From f161d633bbb8d82df9df3c947bbaab3503d79e75 Mon Sep 17 00:00:00 2001 From: andyhuangdev Date: Sat, 22 Aug 2026 11:36:04 +0800 Subject: [PATCH 1/2] HDDS-11470. Report OM snapshot install failures to Ratis --- .../om/ratis/OzoneManagerStateMachine.java | 7 ++++- .../ratis/TestOzoneManagerStateMachine.java | 30 +++++++++++++++++++ 2 files changed, 36 insertions(+), 1 deletion(-) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java index feeda4ca72be..e34978632463 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java @@ -634,7 +634,12 @@ public CompletableFuture notifyInstallSnapshotFromLeader( return CompletableFuture.supplyAsync( () -> { try { - return ozoneManager.installSnapshotFromLeader(leaderNodeId); + TermIndex termIndex = ozoneManager.installSnapshotFromLeader(leaderNodeId); + if (termIndex == null) { + throw new CompletionException( + new IOException("Failed to install snapshot from OM leader " + leaderNodeId)); + } + return termIndex; } catch (IOException e) { throw new CompletionException(e); } diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java index b9fb82786e13..b3fc3e4f41c2 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java @@ -829,6 +829,27 @@ public void testNotifyConfigurationChanged() { // --- notifySnapshotInstalled tests --- + @Test + public void testNotifyInstallSnapshotFromLeaderSuccess() throws Exception { + TermIndex termIndex = TermIndex.valueOf(1, 10); + when(om.installSnapshotFromLeader("leader-om")).thenReturn(termIndex); + + CompletableFuture future = sm.notifyInstallSnapshotFromLeader( + roleInfoWithLeader("leader-om"), TermIndex.valueOf(1, 11)); + + assertEquals(termIndex, future.get()); + } + + @Test + public void testNotifyInstallSnapshotFromLeaderFailure() { + CompletableFuture future = sm.notifyInstallSnapshotFromLeader( + roleInfoWithLeader("leader-om"), TermIndex.valueOf(1, 11)); + + ExecutionException exception = assertThrows(ExecutionException.class, future::get); + + assertInstanceOf(IOException.class, exception.getCause()); + } + @Test public void testNotifySnapshotInstalledSuccess() { RaftPeer localPeer = RaftPeer.newBuilder() @@ -1044,6 +1065,15 @@ private RaftClientRequest buildClientRequest( .build(); } + private RaftProtos.RoleInfoProto roleInfoWithLeader(String leaderId) { + return RaftProtos.RoleInfoProto.newBuilder() + .setFollowerInfo(RaftProtos.FollowerInfoProto.newBuilder() + .setLeaderInfo(RaftProtos.ServerRpcProto.newBuilder() + .setId(RaftProtos.RaftPeerProto.newBuilder() + .setId(ByteString.copyFromUtf8(leaderId))))) + .build(); + } + @Test public void testRatisEventsRecording() { OzoneConfiguration conf = new OzoneConfiguration(); From e6e998adbc4b701aa87d055c5b5fc8f2e9bd4960 Mon Sep 17 00:00:00 2001 From: andyhuangdev Date: Mon, 24 Aug 2026 15:56:29 +0800 Subject: [PATCH 2/2] HDDS-11470. Propagate OM checkpoint installation failures --- .../hadoop/ozone/om/TestOMRatisSnapshots.java | 8 ++++--- .../apache/hadoop/ozone/om/OzoneManager.java | 24 ++++++++++++++----- .../om/ratis/OzoneManagerStateMachine.java | 7 +----- .../ratis/TestOzoneManagerStateMachine.java | 15 ++++++++++-- 4 files changed, 37 insertions(+), 17 deletions(-) diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestOMRatisSnapshots.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestOMRatisSnapshots.java index 1bfab3c3c6bc..5a5739431e6c 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestOMRatisSnapshots.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestOMRatisSnapshots.java @@ -29,6 +29,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import java.io.File; @@ -626,9 +627,10 @@ public void testInstallSnapshotFailedBackupRestoresDbDir() throws Exception { "Simulated backup move failure for test")); followerOM.setExitManagerForTesting(new DummyExitManager()); try { - TermIndex termIndex = followerOM.installCheckpoint( - leaderOMNodeId, leaderCheckpointLocation, leaderCheckpointTrxnInfo); - assertNull(termIndex, "Install should have been reported as failed"); + IOException exception = assertThrows(IOException.class, () -> + followerOM.installCheckpoint(leaderOMNodeId, leaderCheckpointLocation, + leaderCheckpointTrxnInfo)); + assertThat(exception).hasMessageContaining("Cannot replace DB"); // Everything present before the aborted install must still be present. assertThat(topLevelNames(followerMetaDir)).containsAll(namesBefore); 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 4254a634b3d4..08855a784c64 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 @@ -4260,7 +4260,7 @@ public synchronized TermIndex installSnapshotFromLeader(String leaderId) throws } termIndex = installCheckpoint(leaderId, checkpointLocation); } catch (Exception ex) { - LOG.error("Failed to install snapshot from Leader OM.", ex); + throw new IOException("Failed to install snapshot from Leader " + leaderId, ex); } finally { cleanupCheckpoint(omDBCheckpoint); } @@ -4317,6 +4317,7 @@ TermIndex installCheckpoint(String leaderId, Path checkpointLocation, long startTime = Time.monotonicNow(); File oldDBLocation = metadataManager.getStore().getDbLocation(); Path omDbPath = Paths.get(checkpointLocation.toString(), OM_DB_NAME); + IOException installFailure = null; try { // Stop Background services keyManager.stop(); @@ -4328,20 +4329,16 @@ TermIndex installCheckpoint(String leaderId, Path checkpointLocation, // pending transactions in the buffer, they are discarded. omRatisServer.getOmStateMachine().pause(); } catch (Exception e) { - LOG.error("Failed to stop/ pause the services. Cannot proceed with " + - "installing the new checkpoint."); // Stop the checkpoint install process and restart the services. keyManager.start(configuration); startSecretManagerIfNecessary(); startTrashEmptier(configuration); - throw e; + throw newInstallCheckpointException(checkpointTrxnInfo, "stop/pause services", e); } - File dbBackup = null; TermIndex termIndex = omRatisServer.getLastAppliedTermIndex(); long term = termIndex.getTerm(); long lastAppliedIndex = termIndex.getIndex(); - // Check if current applied log index is smaller than the downloaded // checkpoint transaction index. If yes, proceed by stopping the ratis // server so that the OM state can be re-initialized. If no then do not @@ -4383,6 +4380,7 @@ TermIndex installCheckpoint(String leaderId, Path checkpointLocation, LOG.error("Failed to install Snapshot from {} as OM failed to replace" + " DB with downloaded checkpoint. Reloading old OM state.", leaderId, e); + installFailure = newInstallCheckpointException(checkpointTrxnInfo, "replace DB", e); } } else { LOG.warn("Cannot proceed with InstallSnapshot as OM is at TermIndex {} " + @@ -4449,6 +4447,8 @@ TermIndex installCheckpoint(String leaderId, Path checkpointLocation, dbBackup, e); } + throwIfInstallCheckpointFailed(installFailure); + if (lastAppliedIndex != checkpointTrxnInfo.getTransactionIndex()) { // Install Snapshot failed and old state was reloaded. Return null to // Ratis to indicate that installation failed. @@ -4464,6 +4464,18 @@ TermIndex installCheckpoint(String leaderId, Path checkpointLocation, return newTermIndex; } + private static IOException newInstallCheckpointException(TransactionInfo checkpointTrxnInfo, + String operation, Exception cause) { + return new IOException("Failed to install checkpoint " + checkpointTrxnInfo + + ": Cannot " + operation + '.', cause); + } + + private static void throwIfInstallCheckpointFailed(IOException installFailure) throws IOException { + if (installFailure != null) { + throw installFailure; + } + } + private void buildDBCheckpointInstallAuditLog(String leaderId, long term, long lastAppliedIndex) { Map auditMap = new LinkedHashMap<>(); auditMap.put(AUDIT_PARAM_LEADER_ID, leaderId); diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java index e34978632463..feeda4ca72be 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ratis/OzoneManagerStateMachine.java @@ -634,12 +634,7 @@ public CompletableFuture notifyInstallSnapshotFromLeader( return CompletableFuture.supplyAsync( () -> { try { - TermIndex termIndex = ozoneManager.installSnapshotFromLeader(leaderNodeId); - if (termIndex == null) { - throw new CompletionException( - new IOException("Failed to install snapshot from OM leader " + leaderNodeId)); - } - return termIndex; + return ozoneManager.installSnapshotFromLeader(leaderNodeId); } catch (IOException e) { throw new CompletionException(e); } diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java index b3fc3e4f41c2..34f952ab3c35 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ratis/TestOzoneManagerStateMachine.java @@ -841,13 +841,24 @@ public void testNotifyInstallSnapshotFromLeaderSuccess() throws Exception { } @Test - public void testNotifyInstallSnapshotFromLeaderFailure() { + public void testNotifyInstallSnapshotFromLeaderNullResult() throws Exception { + CompletableFuture future = sm.notifyInstallSnapshotFromLeader( + roleInfoWithLeader("leader-om"), TermIndex.valueOf(1, 11)); + + assertNull(future.get()); + } + + @Test + public void testNotifyInstallSnapshotFromLeaderFailure() throws Exception { + IOException failure = new IOException("Failed to install checkpoint"); + doThrow(failure).when(om).installSnapshotFromLeader("leader-om"); + CompletableFuture future = sm.notifyInstallSnapshotFromLeader( roleInfoWithLeader("leader-om"), TermIndex.valueOf(1, 11)); ExecutionException exception = assertThrows(ExecutionException.class, future::get); - assertInstanceOf(IOException.class, exception.getCause()); + assertSame(failure, exception.getCause()); } @Test