diff --git a/hadoop-hdds/interface-client/src/main/proto/hdds.proto b/hadoop-hdds/interface-client/src/main/proto/hdds.proto index 53ad3db88221..fe19f300089d 100644 --- a/hadoop-hdds/interface-client/src/main/proto/hdds.proto +++ b/hadoop-hdds/interface-client/src/main/proto/hdds.proto @@ -417,6 +417,7 @@ message BlockID { optional uint64 blockCommitSequenceId = 2 [default = 0]; } +// Deprecated, please use UpgradeStatus instead; kept for compatibility with old clients. message UpgradeFinalizationStatus { enum Status { ALREADY_FINALIZED = 1; @@ -430,10 +431,13 @@ message UpgradeFinalizationStatus { } message UpgradeStatus { - optional bool scmFinalized = 1; - optional int32 numDatanodesFinalized = 2; - optional int32 numDatanodesTotal = 3; - optional bool shouldFinalize = 4; + optional bool hddsFinalized = 1; + optional bool scmFinalized = 2; + optional int32 numDatanodesFinalized = 3; + optional int32 numDatanodesTotal = 4; + optional uint32 scmApparentVersion = 5; + optional uint32 minDatanodeApparentVersion = 6; + optional uint32 maxDatanodeApparentVersion = 7; } /** diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java index 96eafa0b2aa8..a3afe1f2506d 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java @@ -152,6 +152,8 @@ default int getAllNodeCount() { default DatanodeFinalizationCounts getDatanodeFinalizationCounts() { int finalizedNodes = 0; int totalHealthyNodes = 0; + int minApparentVersion = Integer.MAX_VALUE; + int maxApparentVersion = 0; for (DatanodeInfo dn : getAllNodes()) { try { @@ -173,6 +175,10 @@ default DatanodeFinalizationCounts getDatanodeFinalizationCounts() { ComponentVersion dnApparentVersion = dn.getLastKnownApparentVersion(); ComponentVersion dnSoftwareVersion = dn.getLastKnownSoftwareVersion(); + int dnApparentVersionInt = dnApparentVersion.serialize(); + minApparentVersion = Math.min(minApparentVersion, dnApparentVersionInt); + maxApparentVersion = Math.max(maxApparentVersion, dnApparentVersionInt); + if (!dnApparentVersion.equals(dnSoftwareVersion)) { // Datanode has not yet finalized LOG.debug("Datanode {} has not yet finalized: apparent version={}, software version={}", @@ -187,7 +193,17 @@ default DatanodeFinalizationCounts getDatanodeFinalizationCounts() { } } - return new DatanodeFinalizationCounts(finalizedNodes, totalHealthyNodes); + if (minApparentVersion == Integer.MAX_VALUE) { + // No healthy datanode with version info was found + minApparentVersion = 0; + } + + return DatanodeFinalizationCounts.newBuilder() + .setNumFinalizedDatanodes(finalizedNodes) + .setTotalHealthyDatanodes(totalHealthyNodes) + .setMinApparentVersion(minApparentVersion) + .setMaxApparentVersion(maxApparentVersion) + .build(); } /** @@ -498,11 +514,18 @@ default void removeNode(DatanodeDetails datanodeDetails) throws NodeNotFoundExce final class DatanodeFinalizationCounts { private final int numFinalizedDatanodes; private final int totalHealthyDatanodes; + private final int minApparentVersion; + private final int maxApparentVersion; + + private DatanodeFinalizationCounts(Builder b) { + this.numFinalizedDatanodes = b.numFinalizedDatanodes; + this.totalHealthyDatanodes = b.totalHealthyDatanodes; + this.minApparentVersion = b.minApparentVersion; + this.maxApparentVersion = b.maxApparentVersion; + } - public DatanodeFinalizationCounts(int numFinalizedDatanodes, - int totalHealthyDatanodes) { - this.numFinalizedDatanodes = numFinalizedDatanodes; - this.totalHealthyDatanodes = totalHealthyDatanodes; + public static Builder newBuilder() { + return new Builder(); } public int getNumFinalizedDatanodes() { @@ -516,5 +539,47 @@ public int getTotalHealthyDatanodes() { public boolean allNodesFinalized() { return numFinalizedDatanodes == totalHealthyDatanodes; } + + public int getMinApparentVersion() { + return minApparentVersion; + } + + public int getMaxApparentVersion() { + return maxApparentVersion; + } + + /** + * Builder for {@link DatanodeFinalizationCounts}. + */ + public static final class Builder { + private int numFinalizedDatanodes; + private int totalHealthyDatanodes; + private int minApparentVersion; + private int maxApparentVersion; + + public Builder setNumFinalizedDatanodes(int value) { + this.numFinalizedDatanodes = value; + return this; + } + + public Builder setTotalHealthyDatanodes(int value) { + this.totalHealthyDatanodes = value; + return this; + } + + public Builder setMinApparentVersion(int value) { + this.minApparentVersion = value; + return this; + } + + public Builder setMaxApparentVersion(int value) { + this.maxApparentVersion = value; + return this; + } + + public DatanodeFinalizationCounts build() { + return new DatanodeFinalizationCounts(this); + } + } } } 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 ac5f17a5c6e2..623016457006 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 @@ -1214,18 +1214,26 @@ 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); + } + boolean scmFinalized = !scm.getVersionManager().needsFinalization(); NodeManager.DatanodeFinalizationCounts datanodeFinalizationCounts = scm.getScmNodeManager().getDatanodeFinalizationCounts(); int finalizedDatanodes = datanodeFinalizationCounts.getNumFinalizedDatanodes(); int healthyDatanodes = datanodeFinalizationCounts.getTotalHealthyDatanodes(); - boolean shouldFinalize = scmFinalized && datanodeFinalizationCounts.allNodesFinalized() && !scm.isInSafeMode(); + boolean hddsFinalized = scmFinalized && datanodeFinalizationCounts.allNodesFinalized(); HddsProtos.UpgradeStatus result = HddsProtos.UpgradeStatus.newBuilder() .setScmFinalized(scmFinalized) .setNumDatanodesFinalized(finalizedDatanodes) .setNumDatanodesTotal(healthyDatanodes) - .setShouldFinalize(shouldFinalize) + .setHddsFinalized(hddsFinalized) + .setScmApparentVersion(scm.getVersionManager().getApparentVersion().serialize()) + .setMinDatanodeApparentVersion(datanodeFinalizationCounts.getMinApparentVersion()) + .setMaxDatanodeApparentVersion(datanodeFinalizationCounts.getMaxApparentVersion()) .build(); AUDIT.logReadSuccess(buildAuditMessageForSuccess(SCMAction.QUERY_UPGRADE_STATUS, null)); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java index 8b0ea537d4e4..b7b7c8bdf8ea 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java @@ -684,6 +684,36 @@ public void testDatanodeFinalizedCounterTracksRegistrationAndRemoveNode() } } + @Test + public void testDatanodeFinalizationCountsTracksApparentVersionRange() + throws IOException, AuthenticationException { + try (SCMNodeManager nodeManager = createNodeManager(getConf())) { + // Two healthy datanodes reporting different apparent versions. + registerWithCapacity(nodeManager, + toVersionProto(HDDSLayoutFeature.INITIAL_VERSION, HDDSVersion.SOFTWARE_VERSION), success); + registerWithCapacity(nodeManager, defaultVersionProto(), success); + + NodeManager.DatanodeFinalizationCounts counts = nodeManager.getDatanodeFinalizationCounts(); + assertEquals(2, counts.getTotalHealthyDatanodes()); + assertEquals(HDDSLayoutFeature.INITIAL_VERSION.serialize(), counts.getMinApparentVersion(), + "Min apparent version should reflect the least-finalized datanode"); + assertEquals(HDDSVersion.SOFTWARE_VERSION.serialize(), counts.getMaxApparentVersion(), + "Max apparent version should reflect the most-finalized datanode"); + } + } + + @Test + public void testDatanodeFinalizationCountsWithNoHealthyDatanodes() + throws IOException, AuthenticationException { + try (SCMNodeManager nodeManager = createNodeManager(getConf())) { + NodeManager.DatanodeFinalizationCounts counts = nodeManager.getDatanodeFinalizationCounts(); + assertEquals(0, counts.getTotalHealthyDatanodes()); + // The MAX_VALUE sentinel used to find the min must be normalized to 0 when there is nothing to count. + assertEquals(0, counts.getMinApparentVersion()); + assertEquals(0, counts.getMaxApparentVersion()); + } + } + private static Stream ineligibleHealthStates() { return Stream.of( Arguments.of(NodeStatus.inServiceStale()), 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 61ec18c3e1cc..8ac74e4b86bd 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 @@ -21,7 +21,6 @@ import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_READONLY_ADMINISTRATORS; import static org.apache.hadoop.ozone.upgrade.UpgradeFinalization.Status.ALREADY_FINALIZED; import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; @@ -50,12 +49,14 @@ 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.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.SCMNodeDetails; import org.apache.hadoop.hdds.scm.pipeline.PipelineID; 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.ProtocolMessageMetrics; @@ -284,24 +285,23 @@ public void testQueryUpgradeStatus() throws Exception { // No datanodes registered assertEquals(0, status.getNumDatanodesFinalized()); assertEquals(0, status.getNumDatanodesTotal()); - assertTrue(status.getShouldFinalize()); + assertTrue(status.getHddsFinalized()); } @Test - public void testQueryUpgradeStatusInSafemode() throws Exception { - // mockSafeModeManager defaults to returning true for getInSafeMode() - when(mockSafeModeManager.getInSafeMode()).thenReturn(true); - assertTrue(scm.isInSafeMode()); - - HddsProtos.UpgradeStatus status = server.queryUpgradeStatus(); + public void testQueryUpgradeStatusInSafemode() { + // Put SCM into safe mode via the context the server consults. + scm.getScmContext().updateSafeModeStatus(SafeModeStatus.INITIAL); + try { + assertTrue(scm.getScmContext().isInSafeMode()); - // SCM starts already finalized in tests - assertTrue(status.getScmFinalized()); - // No datanodes registered - assertEquals(0, status.getNumDatanodesFinalized()); - assertEquals(0, status.getNumDatanodesTotal()); - // shouldFinalize is false because SCM is in safe mode - assertFalse(status.getShouldFinalize()); + // 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, 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 411955951df6..232f3347e9e3 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 @@ -17,13 +17,16 @@ package org.apache.hadoop.ozone.admin.upgrade; +import com.google.common.annotations.VisibleForTesting; import java.util.concurrent.Callable; +import java.util.concurrent.TimeUnit; import org.apache.hadoop.hdds.cli.AbstractSubcommand; import org.apache.hadoop.hdds.cli.HddsVersionProvider; import org.apache.hadoop.ozone.OzoneManagerVersion; import org.apache.hadoop.ozone.admin.om.OmAddressOptions; import org.apache.hadoop.ozone.client.rpc.RpcClient; import org.apache.hadoop.ozone.om.protocol.OzoneManagerProtocol; +import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.QueryUpgradeStatusResponse; import picocli.CommandLine; /** @@ -31,15 +34,25 @@ */ @CommandLine.Command( name = "finalize", - description = "Initiates the the process to finalize a cluster upgrade.", + description = "Initiates the process to finalize a cluster upgrade. This command is idempotent.", mixinStandardHelpOptions = true, versionProvider = HddsVersionProvider.class ) public class FinalizeSubCommand extends AbstractSubcommand implements Callable { + /** Poll cadence used by {@code --wait}. Overridable from tests via {@link #setPollIntervalMillis(long)}. */ + private long pollIntervalMillis = TimeUnit.SECONDS.toMillis(5); + @CommandLine.Mixin private OmAddressOptions.OptionalServiceIdOrHostMixin omAddressOptions; + @CommandLine.Option(names = {"--wait"}, + defaultValue = "false", + description = "After initiating finalization, poll the cluster status until the entire cluster (OM, SCM, " + + "and all healthy datanodes) is finalized. Interrupt with Ctrl-C to stop waiting; finalization " + + "continues on the server.") + private boolean wait; + @Override public Integer call() throws Exception { try (OzoneManagerProtocol client = getClient()) { @@ -50,12 +63,63 @@ public Integer call() throws Exception { return 1; } 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)."); + return waitForFinalization(client); + } out().println("Cluster finalization has been started. Monitor progress with `ozone admin upgrade status`"); } return 0; } + /** + * Polls the cluster status until OM, SCM and all healthy datanodes report finalized, or the operator + * interrupts the command. A failed status query is reported to stderr and retried on the next poll, + * so transient RPC errors do not abort the wait. {@code --wait} is safe to re-run: if the cluster is + * already finalized, the first poll returns done and the command exits 0. + */ + private int waitForFinalization(OzoneManagerProtocol client) { + while (true) { + QueryUpgradeStatusResponse status = null; + try { + status = client.queryUpgradeStatus(); + } catch (Exception e) { + err().println("Failed to query upgrade status: " + e.getMessage() + + ". Retrying, or use `ozone admin upgrade status` to monitor progress."); + } + + if (status != null) { + // Finalization checks before sleeping, so an already-finalized cluster returns without waiting. + if (status.getClusterFinalized()) { + out().println("Finalization complete."); + return 0; + } + + if (isVerbose()) { + StatusSubCommand.printVerbose(status, out()); + } else { + StatusSubCommand.printBasic(status, out()); + } + out().flush(); + } + + try { + Thread.sleep(pollIntervalMillis); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + out().println("Waiting interrupted. Use `ozone admin upgrade status` to monitor progress."); + return 1; + } + } + } + protected OzoneManagerProtocol getClient() throws Exception { return omAddressOptions.newClient(); } + + @VisibleForTesting + void setPollIntervalMillis(long millis) { + this.pollIntervalMillis = millis; + } } diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/ozone/admin/upgrade/StatusSubCommand.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/ozone/admin/upgrade/StatusSubCommand.java index bc3a18cd1b94..fa1ef06bb586 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/ozone/admin/upgrade/StatusSubCommand.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/ozone/admin/upgrade/StatusSubCommand.java @@ -17,14 +17,18 @@ package org.apache.hadoop.ozone.admin.upgrade; +import java.io.PrintWriter; import java.util.concurrent.Callable; +import org.apache.hadoop.hdds.HDDSVersion; import org.apache.hadoop.hdds.cli.AbstractSubcommand; import org.apache.hadoop.hdds.cli.HddsVersionProvider; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos; +import org.apache.hadoop.hdds.server.JsonUtils; import org.apache.hadoop.ozone.OzoneManagerVersion; import org.apache.hadoop.ozone.admin.om.OmAddressOptions; import org.apache.hadoop.ozone.client.rpc.RpcClient; import org.apache.hadoop.ozone.om.protocol.OzoneManagerProtocol; -import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos; +import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.QueryUpgradeStatusResponse; import picocli.CommandLine; /** @@ -42,6 +46,11 @@ public class StatusSubCommand extends AbstractSubcommand implements Callable thrown = new AtomicReference<>(); + Thread runner = new Thread(() -> { + try { + result.set(cmd.call()); + } catch (Exception e) { + // The command handles interruption internally and should not throw. + thrown.set(e); + } + }); + runner.start(); + // Give the runner a moment to complete the first poll and start sleeping. + Thread.sleep(50); + runner.interrupt(); + runner.join(5_000); + + assertNull(thrown.get(), "wait command should not throw on interrupt"); + String output = outContent.toString(DEFAULT_ENCODING); + assertTrue(output.contains("Waiting interrupted")); + // With check-before-sleep, the first poll runs before the interrupting sleep. + verify(omClient, times(1)).queryUpgradeStatus(); + // The command swallows the interrupt, but reports non-zero since finalization was not confirmed. + assertEquals(1, result.get()); + } + + @Test + public void testWaitFlagIsResumableAfterCancel() throws Exception { + // First invocation: one in-progress poll, then blocked in Thread.sleep so we can interrupt it. + cmd.setPollIntervalMillis(60_000); + // First invocation's poll: in progress (so it sleeps); second invocation's poll: finalized. + when(omClient.queryUpgradeStatus()) + .thenReturn(inProgressStatus(0, 3)) + .thenReturn(finalizedStatus(3, 3)); + + new CommandLine(cmd).parseArgs("--wait"); + + AtomicInteger firstResult = new AtomicInteger(-1); + AtomicReference thrown = new AtomicReference<>(); + Thread runner = new Thread(() -> { + try { + firstResult.set(cmd.call()); + } catch (Exception e) { + // First call should swallow the interrupt cleanly and not throw. + thrown.set(e); + } + }); + runner.start(); + Thread.sleep(50); + runner.interrupt(); + runner.join(5_000); + + assertNull(thrown.get(), "first wait command should not throw on interrupt"); + String firstOutput = outContent.toString(DEFAULT_ENCODING); + assertTrue(firstOutput.contains("Cluster finalization has been started")); + assertTrue(firstOutput.contains("Waiting interrupted")); + assertEquals(1, firstResult.get()); + // First invocation performed exactly one poll before being interrupted during the sleep. + verify(omClient, times(1)).queryUpgradeStatus(); + + // Second invocation — reset output capture and shrink the poll interval so it exits promptly. + outContent.reset(); + cmd.setPollIntervalMillis(1); + + new CommandLine(cmd).parseArgs("--wait"); + assertEquals(0, cmd.call()); + + String secondOutput = outContent.toString(DEFAULT_ENCODING); + assertTrue(secondOutput.contains("Cluster finalization has been started")); + assertTrue(secondOutput.contains("Finalization complete.")); + // finalizeUpgrade must have been issued on both runs (idempotent server-side). + verify(omClient, times(2)).finalizeUpgrade(); + // queryUpgradeStatus ran once per invocation. + verify(omClient, times(2)).queryUpgradeStatus(); + // The client is closed after each try-with-resources block. + verify(omClient, times(2)).close(); + } + + @Test + public void testWaitFlagWithVerbosePrintsFullStatus() throws Exception { + verbose = true; + // Return one in-progress poll first so the verbose status is printed. + when(omClient.queryUpgradeStatus()) + .thenReturn(inProgressStatus(1, 2)) + .thenReturn(finalizedStatus(2, 2)); + + new CommandLine(cmd).parseArgs("--wait"); + assertEquals(0, cmd.call()); + + String output = outContent.toString(DEFAULT_ENCODING); + assertTrue(output.contains("OM Finalized?")); + assertTrue(output.contains("SCM Finalized?")); + assertTrue(output.contains("OM Apparent Version:")); + assertTrue(output.contains("SCM Apparent Version:")); + assertTrue(output.contains("Min Datanode Apparent Version:")); + assertFalse(output.contains("Waiting for finalization:")); + assertTrue(output.contains("Finalization complete.")); + } + + private static QueryUpgradeStatusResponse inProgressStatus(int dnFinalized, int dnTotal) { + return QueryUpgradeStatusResponse.newBuilder() + .setOmFinalized(false) + .setClusterFinalized(false) + .setHddsStatus(HddsProtos.UpgradeStatus.newBuilder() + .setScmFinalized(false) + .setNumDatanodesFinalized(dnFinalized) + .setNumDatanodesTotal(dnTotal) + .build()) + .build(); + } + + private static QueryUpgradeStatusResponse finalizedStatus(int dnFinalized, int dnTotal) { + return QueryUpgradeStatusResponse.newBuilder() + .setOmFinalized(true) + .setClusterFinalized(true) + .setHddsStatus(HddsProtos.UpgradeStatus.newBuilder() + .setScmFinalized(true) + .setNumDatanodesFinalized(dnFinalized) + .setNumDatanodesTotal(dnTotal) + .build()) + .build(); + } + private ServiceInfoEx serviceInfoWithVersion(OzoneManagerVersion version) { ServiceInfo serviceInfo = new ServiceInfo.Builder() .setNodeType(HddsProtos.NodeType.OM) diff --git a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/ozone/admin/upgrade/TestStatusSubCommand.java b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/ozone/admin/upgrade/TestStatusSubCommand.java index 26cfef1ffe63..76ff43998e18 100644 --- a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/ozone/admin/upgrade/TestStatusSubCommand.java +++ b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/ozone/admin/upgrade/TestStatusSubCommand.java @@ -18,6 +18,7 @@ package org.apache.hadoop.ozone.admin.upgrade; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; @@ -25,11 +26,14 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.PrintStream; import java.nio.charset.StandardCharsets; import java.util.Collections; +import org.apache.hadoop.hdds.HDDSVersion; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.ozone.OzoneManagerVersion; import org.apache.hadoop.ozone.om.helpers.ServiceInfo; @@ -47,6 +51,7 @@ public class TestStatusSubCommand { private static final String DEFAULT_ENCODING = StandardCharsets.UTF_8.name(); + private static final ObjectMapper JSON = new ObjectMapper(); private final ByteArrayOutputStream outContent = new ByteArrayOutputStream(); private final ByteArrayOutputStream errContent = new ByteArrayOutputStream(); @@ -54,17 +59,24 @@ public class TestStatusSubCommand { private final PrintStream originalErr = System.err; private OzoneManagerProtocol omClient; private StatusSubCommand cmd; + private boolean verbose; @BeforeEach public void setup() throws IOException { omClient = mock(OzoneManagerProtocol.class); when(omClient.getServiceInfo()).thenReturn(serviceInfoWithVersion(OzoneManagerVersion.ZDU)); + verbose = false; cmd = new StatusSubCommand() { @Override protected OzoneManagerProtocol getClient() throws Exception { return omClient; } + + @Override + protected boolean isVerbose() { + return verbose; + } }; System.setOut(new PrintStream(outContent, false, DEFAULT_ENCODING)); System.setErr(new PrintStream(errContent, false, DEFAULT_ENCODING)); @@ -82,7 +94,7 @@ public void testStatusCommandPrintsUpgradeStatus() throws Exception { .setScmFinalized(false) .setNumDatanodesFinalized(1) .setNumDatanodesTotal(3) - .setShouldFinalize(true) + .setHddsFinalized(true) .build(); OzoneManagerProtocolProtos.QueryUpgradeStatusResponse response = @@ -100,6 +112,8 @@ public void testStatusCommandPrintsUpgradeStatus() throws Exception { assertTrue(output.contains("OM Finalized? false")); assertTrue(output.contains("SCM Finalized? false")); assertTrue(output.contains("Datanodes finalized: 1/3")); + // Without --verbose the apparent versions are not shown. + assertFalse(output.contains("Apparent Version")); verify(omClient).queryUpgradeStatus(); } @@ -122,6 +136,82 @@ public void testNonZduServerPrintsErrorAndReturnsNonZero() throws Exception { verify(omClient, never()).queryUpgradeStatus(); } + @Test + public void testJsonOutput() throws Exception { + int omVersion = OzoneManagerVersion.ZDU.serialize(); + int hddsVersion = HDDSVersion.SOFTWARE_VERSION.serialize(); + OzoneManagerProtocolProtos.QueryUpgradeStatusResponse response = + OzoneManagerProtocolProtos.QueryUpgradeStatusResponse.newBuilder() + .setOmFinalized(true) + .setOmApparentVersion(omVersion) + .setHddsStatus(HddsProtos.UpgradeStatus.newBuilder() + .setScmFinalized(true) + .setNumDatanodesFinalized(2) + .setNumDatanodesTotal(3) + .setScmApparentVersion(hddsVersion) + .setMinDatanodeApparentVersion(hddsVersion) + .setMaxDatanodeApparentVersion(hddsVersion) + .build()) + .build(); + when(omClient.queryUpgradeStatus()).thenReturn(response); + + // JSON output includes every field, including the apparent versions. + new CommandLine(cmd).parseArgs("--json"); + assertEquals(0, cmd.call()); + String jsonOutput = outContent.toString(DEFAULT_ENCODING); + + JsonNode root = JSON.readTree(jsonOutput); + assertTrue(root.path("omFinalized").asBoolean()); + assertTrue(root.path("scmFinalized").asBoolean()); + assertEquals(2, root.path("datanodesFinalized").asInt()); + assertEquals(3, root.path("datanodesTotal").asInt()); + assertEquals(OzoneManagerVersion.ZDU.toString(), root.path("omApparentVersion").asText()); + assertEquals(HDDSVersion.SOFTWARE_VERSION.toString(), root.path("scmApparentVersion").asText()); + assertEquals(HDDSVersion.SOFTWARE_VERSION.toString(), root.path("minDatanodeApparentVersion").asText()); + assertEquals(HDDSVersion.SOFTWARE_VERSION.toString(), root.path("maxDatanodeApparentVersion").asText()); + + // The --verbose flag only affects the human-readable output; JSON output is identical. + outContent.reset(); + verbose = true; + new CommandLine(cmd).parseArgs("--json"); + assertEquals(0, cmd.call()); + assertEquals(jsonOutput, outContent.toString(DEFAULT_ENCODING)); + } + + @Test + public void testVerboseTextOutputIncludesVersions() throws Exception { + verbose = true; + int omVersion = OzoneManagerVersion.ZDU.serialize(); + int hddsVersion = HDDSVersion.SOFTWARE_VERSION.serialize(); + OzoneManagerProtocolProtos.QueryUpgradeStatusResponse response = + OzoneManagerProtocolProtos.QueryUpgradeStatusResponse.newBuilder() + .setOmFinalized(true) + .setOmApparentVersion(omVersion) + .setHddsStatus(HddsProtos.UpgradeStatus.newBuilder() + .setScmFinalized(true) + .setNumDatanodesFinalized(3) + .setNumDatanodesTotal(3) + .setScmApparentVersion(hddsVersion) + .setMinDatanodeApparentVersion(hddsVersion) + .setMaxDatanodeApparentVersion(hddsVersion) + .build()) + .build(); + when(omClient.queryUpgradeStatus()).thenReturn(response); + + new CommandLine(cmd).parseArgs(); + assertEquals(0, cmd.call()); + + String output = outContent.toString(DEFAULT_ENCODING); + assertTrue(output.contains("OM Finalized?")); + assertTrue(output.contains("OM Apparent Version:")); + assertTrue(output.contains("SCM Finalized?")); + assertTrue(output.contains("SCM Apparent Version:")); + assertTrue(output.contains("Min Datanode Apparent Version:")); + assertTrue(output.contains("Max Datanode Apparent Version:")); + assertTrue(output.contains(OzoneManagerVersion.ZDU.toString())); + assertTrue(output.contains(HDDSVersion.SOFTWARE_VERSION.toString())); + } + private ServiceInfoEx serviceInfoWithVersion(OzoneManagerVersion version) { ServiceInfo serviceInfo = new ServiceInfo.Builder() .setNodeType(HddsProtos.NodeType.OM) 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 3f735828b4c5..04defd997745 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 @@ -51,7 +51,7 @@ public static void waitForFinalizationFromClient(StorageContainerLocationProtoco 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); - return status.getShouldFinalize(); + return status.getHddsFinalized(); }); } diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/OMUpgradeTestUtils.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/OMUpgradeTestUtils.java index cd2d704eb462..675bfdcac67a 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/OMUpgradeTestUtils.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/OMUpgradeTestUtils.java @@ -17,20 +17,24 @@ package org.apache.hadoop.ozone.om; -import static org.apache.hadoop.ozone.upgrade.UpgradeFinalization.Status.FINALIZATION_DONE; import static org.apache.ozone.test.GenericTestUtils.waitFor; import static org.junit.jupiter.api.Assertions.fail; import java.io.IOException; import java.util.concurrent.TimeoutException; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.ozone.om.protocol.OzoneManagerProtocol; -import org.apache.hadoop.ozone.upgrade.UpgradeFinalization; +import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.QueryUpgradeStatusResponse; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * Utility class to help test OM upgrade scenarios. */ public final class OMUpgradeTestUtils { + private static final Logger LOG = LoggerFactory.getLogger(OMUpgradeTestUtils.class); + private OMUpgradeTestUtils() { // Utility class. } @@ -39,12 +43,14 @@ public static void waitForFinalization(OzoneManagerProtocol omClient) throws TimeoutException, InterruptedException { waitFor(() -> { try { - UpgradeFinalization.StatusAndMessages statusAndMessages = - omClient.queryUpgradeFinalizationProgress("finalize-test", false, - false); - System.out.println("Finalization Messages : " + - statusAndMessages.msgs()); - return statusAndMessages.status().equals(FINALIZATION_DONE); + QueryUpgradeStatusResponse status = omClient.queryUpgradeStatus(); + HddsProtos.UpgradeStatus hdds = status.getHddsStatus(); + LOG.info("Finalization status: omFinalized={}, scmFinalized={}, datanodes={}/{}", + status.getOmFinalized(), hdds.getScmFinalized(), + hdds.getNumDatanodesFinalized(), hdds.getNumDatanodesTotal()); + return status.getOmFinalized() + && hdds.getScmFinalized() + && hdds.getNumDatanodesFinalized() == hdds.getNumDatanodesTotal(); } catch (IOException e) { fail(e.getMessage()); } diff --git a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto index cf60eb31f3af..fec82cdda363 100644 --- a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto +++ b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto @@ -1666,8 +1666,11 @@ message QueryUpgradeStatusRequest { } message QueryUpgradeStatusResponse { - optional bool omFinalized = 1; + // True when OM, SCM and all healthy datanodes are finalized + optional bool clusterFinalized = 1; optional hadoop.hdds.UpgradeStatus hddsStatus = 2; + optional bool omFinalized = 3; + optional uint32 omApparentVersion = 4; } // deprecated, use QueryUpgradeStatusRequest/Response instead. Retained for wire compatibility. 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 5e27bde4a2b1..012e4e857a90 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 @@ -194,6 +194,7 @@ import org.apache.hadoop.hdds.scm.ScmInfo; import org.apache.hadoop.hdds.scm.client.HddsClientUtils; import org.apache.hadoop.hdds.scm.client.ScmTopologyClient; +import org.apache.hadoop.hdds.scm.exceptions.SCMException; import org.apache.hadoop.hdds.scm.ha.SCMNodeInfo; import org.apache.hadoop.hdds.scm.net.NetworkTopology; import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol; @@ -3664,11 +3665,26 @@ public void finalizeUpgrade() throws IOException { @Override public QueryUpgradeStatusResponse queryUpgradeStatus() throws IOException { - HddsProtos.UpgradeStatus scmStatus = scmClient.getContainerClient().queryUpgradeStatus(); + HddsProtos.UpgradeStatus scmStatus; + try { + scmStatus = scmClient.getContainerClient().queryUpgradeStatus(); + } catch (SCMException e) { + // SCM refuses the query while in safe mode + if (e.getResult() == SCMException.ResultCodes.SAFE_MODE_EXCEPTION) { + throw new OMException(e.getMessage(), e, NOT_SUPPORTED_OPERATION); + } + throw e; + } + + boolean omFinalized = !versionManager.needsFinalization(); + boolean hddsFinalized = scmStatus.getHddsFinalized(); + boolean clusterFinalized = omFinalized && hddsFinalized; return QueryUpgradeStatusResponse.newBuilder() - .setOmFinalized(!versionManager.needsFinalization()) + .setOmFinalized(omFinalized) .setHddsStatus(scmStatus) + .setOmApparentVersion(versionManager.getApparentVersion().serialize()) + .setClusterFinalized(clusterFinalized) .build(); } 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 93926ba65944..00cc567bf3d1 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 @@ -23,6 +23,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; +import org.apache.hadoop.hdds.scm.exceptions.SCMException; import org.apache.hadoop.hdds.utils.BackgroundService; import org.apache.hadoop.hdds.utils.BackgroundTask; import org.apache.hadoop.hdds.utils.BackgroundTaskQueue; @@ -116,7 +117,7 @@ public BackgroundTaskResult call() { } HddsProtos.UpgradeStatus upgradeStatus = scmClient.getContainerClient().queryUpgradeStatus(); - if (upgradeStatus.getShouldFinalize()) { + if (upgradeStatus.getHddsFinalized()) { LOG.info("The SCM Upgrade has been finalized. OM will now finalize. Run count {}", run); OzoneManagerProtocolProtos.OMRequest omRequest = OzoneManagerProtocolProtos.OMRequest.newBuilder() @@ -133,8 +134,13 @@ public BackgroundTaskResult call() { LOG.debug("The SCM Upgrade has not been finalized. Run count {}", run); } } catch (Exception e) { - LOG.error("An exception occurred while trying to check the SCM Upgrade status or finalize OM. Run count {}", - run, e); + if (e instanceof SCMException + && ((SCMException) e).getResult() == SCMException.ResultCodes.SAFE_MODE_EXCEPTION) { + LOG.info("SCM is in safe mode; will retry checking the upgrade status on the next run. Run count {}", run); + } else { + LOG.error("An exception occurred while trying to check the SCM Upgrade status or finalize OM. Run count {}", + run, e); + } } } else { LOG.debug("Finalization is not in progress. Run count {}", run); 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 a3aa7be8b5f0..ef0f6ef6eb85 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 @@ -125,7 +125,7 @@ void testNoTasksSubmittedWhenFinalizationNotNeeded() throws Exception { /** * When the OM is the leader, finalization is needed, the finalization command is given and SCM reports - * shouldFinalize=true, a FinalizeUpgrade request should be submitted via Ratis. + * hddsFinalized=true, a FinalizeUpgrade request should be submitted via Ratis. */ @Test void testFinalizationTriggeredWhenScmIsFinalizedAndFinalizationInProgress() throws Exception { @@ -134,7 +134,7 @@ void testFinalizationTriggeredWhenScmIsFinalizedAndFinalizationInProgress() thro HddsProtos.UpgradeStatus scmStatus = HddsProtos.UpgradeStatus.newBuilder() .setScmFinalized(true) - .setShouldFinalize(true) + .setHddsFinalized(true) .setNumDatanodesFinalized(3) .setNumDatanodesTotal(3) .build(); @@ -156,7 +156,7 @@ void testFinalizationTriggeredWhenScmIsFinalizedAndFinalizationInProgress() thro } /** - * When SCM reports shouldFinalize=false (SCM is not yet finalized), + * When SCM reports hddsFinalized=false (SCM is not yet finalized), * no Ratis request should be submitted. */ @Test @@ -166,7 +166,7 @@ void testFinalizationSkippedWhenScmNotYetFinalized() throws Exception { HddsProtos.UpgradeStatus scmStatus = HddsProtos.UpgradeStatus.newBuilder() .setScmFinalized(false) - .setShouldFinalize(false) + .setHddsFinalized(false) .setNumDatanodesFinalized(0) .setNumDatanodesTotal(3) .build(); @@ -248,7 +248,7 @@ void testExceptionFromRatisSubmitIsHandledGracefully() throws Exception { HddsProtos.UpgradeStatus scmStatus = HddsProtos.UpgradeStatus.newBuilder() .setScmFinalized(true) - .setShouldFinalize(true) + .setHddsFinalized(true) .setNumDatanodesFinalized(3) .setNumDatanodesTotal(3) .build(); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java index a2c51a7e88c3..764a82866484 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java @@ -476,7 +476,7 @@ public void testQueryUpgradeStatusDispatch() throws IOException { HddsProtos.UpgradeStatus hddsStatus = HddsProtos.UpgradeStatus.newBuilder() .setScmFinalized(true) - .setShouldFinalize(false) + .setHddsFinalized(false) .setNumDatanodesFinalized(3) .setNumDatanodesTotal(3) .build();