From 40b2f7ba110c24a8ffdd103aeb5761264c069cf7 Mon Sep 17 00:00:00 2001 From: Aryan Gupta Date: Mon, 17 Aug 2026 15:48:07 +0530 Subject: [PATCH 1/5] HDDS-12865. Handle automatic container cache refresh for replica verification. --- .../replicas/BlockExistenceVerifier.java | 6 + .../replicas/BlockVerificationResult.java | 15 ++ .../debug/replicas/ChecksumVerifier.java | 17 ++ .../ozone/debug/replicas/ReplicasVerify.java | 71 +++++-- .../debug/replicas/TestReplicasVerify.java | 187 ++++++++++++++++++ 5 files changed, 284 insertions(+), 12 deletions(-) create mode 100644 hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockExistenceVerifier.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockExistenceVerifier.java index dde79989d02f..5cae34533219 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockExistenceVerifier.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockExistenceVerifier.java @@ -24,6 +24,7 @@ import org.apache.hadoop.hdds.scm.XceiverClientManager; import org.apache.hadoop.hdds.scm.XceiverClientSpi; import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; +import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.hdds.scm.pipeline.Pipeline; import org.apache.hadoop.hdds.scm.storage.ContainerProtocolCalls; import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfo; @@ -67,6 +68,11 @@ public BlockVerificationResult verifyBlock(DatanodeDetails datanode, OmKeyLocati } else { return BlockVerificationResult.failCheck("Block does not exist on this replica"); } + } catch (StorageContainerException e) { + if (e.getResult() == ContainerProtos.Result.NO_SUCH_BLOCK) { + return BlockVerificationResult.failCheckAndRefreshKeyLocation(e.getMessage()); + } + return BlockVerificationResult.failIncomplete(e.getMessage()); } catch (IOException e) { return BlockVerificationResult.failIncomplete(e.getMessage()); } finally { diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockVerificationResult.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockVerificationResult.java index c73635ba4755..e94e45e24973 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockVerificationResult.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/BlockVerificationResult.java @@ -28,11 +28,18 @@ public class BlockVerificationResult { private final boolean completed; private final boolean pass; private final List failures; + private final boolean refreshKeyLocation; public BlockVerificationResult(boolean completed, boolean pass, List failures) { + this(completed, pass, failures, false); + } + + public BlockVerificationResult(boolean completed, boolean pass, List failures, + boolean refreshKeyLocation) { this.completed = completed; this.pass = pass; this.failures = failures; + this.refreshKeyLocation = refreshKeyLocation; } public static BlockVerificationResult pass() { @@ -43,6 +50,10 @@ public static BlockVerificationResult failCheck(String message) { return new BlockVerificationResult(true, false, Collections.singletonList(message)); } + public static BlockVerificationResult failCheckAndRefreshKeyLocation(String message) { + return new BlockVerificationResult(true, false, Collections.singletonList(message), true); + } + public static BlockVerificationResult failIncomplete(String message) { return new BlockVerificationResult(false, false, Collections.singletonList(message)); } @@ -59,4 +70,8 @@ public List getFailures() { return failures; } + public boolean shouldRefreshKeyLocation() { + return refreshKeyLocation; + } + } diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ChecksumVerifier.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ChecksumVerifier.java index 96f8218526d0..c2fb110721f1 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ChecksumVerifier.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ChecksumVerifier.java @@ -26,6 +26,7 @@ import org.apache.hadoop.hdds.scm.OzoneClientConfig; import org.apache.hadoop.hdds.scm.XceiverClientManager; import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; +import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.hdds.scm.pipeline.Pipeline; import org.apache.hadoop.ozone.client.io.BlockInputStreamFactoryImpl; import org.apache.hadoop.ozone.common.OzoneChecksumException; @@ -68,6 +69,11 @@ public BlockVerificationResult verifyBlock(DatanodeDetails datanode, OmKeyLocati return BlockVerificationResult.pass(); } catch (IOException e) { Throwable cause = e.getCause() != null ? e.getCause() : e; + StorageContainerException storageContainerException = findStorageContainerException(e); + if (storageContainerException != null + && storageContainerException.getResult() == org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.NO_SUCH_BLOCK) { + return BlockVerificationResult.failCheckAndRefreshKeyLocation(storageContainerException.getMessage()); + } if (cause instanceof OzoneChecksumException) { return BlockVerificationResult.failCheck(cause.getMessage()); } else { @@ -75,4 +81,15 @@ public BlockVerificationResult verifyBlock(DatanodeDetails datanode, OmKeyLocati } } } + + private static StorageContainerException findStorageContainerException(Throwable throwable) { + Throwable current = throwable; + while (current != null) { + if (current instanceof StorageContainerException) { + return (StorageContainerException) current; + } + current = current.getCause(); + } + return null; + } } diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java index 44a8961f1b4d..fd5f475d4044 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java @@ -341,6 +341,32 @@ void processKey(OzoneClient ozoneClient, String volumeName, String bucketName, S return; } + KeyVerificationResult keyVerificationResult = verifyKey(keyInfo, volumeName, bucketName, keyName); + if (keyVerificationResult.shouldRefreshKeyLocation()) { + OmKeyInfo refreshedKeyInfo = ozoneClient.getProxy().getKeyInfo( + volumeName, bucketName, keyName, true); + keyVerificationResult = verifyKey(refreshedKeyInfo, volumeName, bucketName, keyName); + } + + keyVerificationResult.getKeyNode().put("pass", keyVerificationResult.passed()); + if (keyVerificationResult.passed()) { + keysPassed.incrementAndGet(); + } else { + keysFailed.incrementAndGet(); + allKeysPassed.set(false); + keyVerificationResult.getFailedVerificationTypes().forEach(failedType -> failuresByType + .computeIfAbsent(failedType, k -> new AtomicInteger(0)) + .incrementAndGet() + ); + } + + if (!keyVerificationResult.passed() || allResults) { + keysArray.add(keyVerificationResult.getKeyNode()); + } + } + + private KeyVerificationResult verifyKey( + OmKeyInfo keyInfo, String volumeName, String bucketName, String keyName) { ObjectNode keyNode = JsonUtils.createObjectNode(null); keyNode.put("volumeName", volumeName); keyNode.put("bucketName", bucketName); @@ -348,6 +374,7 @@ void processKey(OzoneClient ozoneClient, String volumeName, String bucketName, S ArrayNode blocksArray = keyNode.putArray("blocks"); boolean keyPass = true; + boolean shouldRefreshKeyLocation = false; Set failedVerificationTypes = new HashSet<>(); for (OmKeyLocationInfo keyLocation : keyInfo.getLatestVersionLocations().getBlocksLatestVersionOnly()) { @@ -378,6 +405,9 @@ void processKey(OzoneClient ozoneClient, String volumeName, String bucketName, S checkNode.put("type", verifier.getType()); checkNode.put("completed", result.isCompleted()); checkNode.put("pass", result.passed()); + if (result.shouldRefreshKeyLocation()) { + shouldRefreshKeyLocation = true; + } ArrayNode failuresArray = checkNode.putArray("failures"); for (String failure : result.getFailures()) { @@ -401,20 +431,37 @@ void processKey(OzoneClient ozoneClient, String volumeName, String bucketName, S } } - keyNode.put("pass", keyPass); - if (keyPass) { - keysPassed.incrementAndGet(); - } else { - keysFailed.incrementAndGet(); - allKeysPassed.set(false); - failedVerificationTypes.forEach(failedType -> failuresByType - .computeIfAbsent(failedType, k -> new AtomicInteger(0)) - .incrementAndGet() - ); + return new KeyVerificationResult(keyNode, keyPass, failedVerificationTypes, shouldRefreshKeyLocation); + } + + private static final class KeyVerificationResult { + private final ObjectNode keyNode; + private final boolean pass; + private final Set failedVerificationTypes; + private final boolean refreshKeyLocation; + + private KeyVerificationResult(ObjectNode keyNode, boolean pass, Set failedVerificationTypes, + boolean refreshKeyLocation) { + this.keyNode = keyNode; + this.pass = pass; + this.failedVerificationTypes = failedVerificationTypes; + this.refreshKeyLocation = refreshKeyLocation; + } + + private ObjectNode getKeyNode() { + return keyNode; + } + + private boolean passed() { + return pass; + } + + private Set getFailedVerificationTypes() { + return failedVerificationTypes; } - if (!keyPass || allResults) { - keysArray.add(keyNode); + private boolean shouldRefreshKeyLocation() { + return refreshKeyLocation; } } diff --git a/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java new file mode 100644 index 000000000000..2e9cf8262d8c --- /dev/null +++ b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java @@ -0,0 +1,187 @@ +/* + * 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.debug.replicas; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.fasterxml.jackson.databind.node.ArrayNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import java.io.IOException; +import java.lang.reflect.Field; +import java.util.Arrays; +import java.util.Collections; +import java.util.Optional; +import java.util.concurrent.atomic.AtomicBoolean; +import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.client.StandaloneReplicationConfig; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos; +import org.apache.hadoop.hdds.scm.pipeline.Pipeline; +import org.apache.hadoop.hdds.scm.pipeline.PipelineID; +import org.apache.hadoop.hdds.server.JsonUtils; +import org.apache.hadoop.ozone.client.OzoneClient; +import org.apache.hadoop.ozone.client.protocol.ClientProtocol; +import org.apache.hadoop.ozone.om.helpers.OmKeyInfo; +import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfo; +import org.apache.hadoop.ozone.om.helpers.OmKeyLocationInfoGroup; +import org.apache.hadoop.ozone.shell.ShellReplicationOptions; +import org.junit.jupiter.api.Test; + +/** + * Unit tests for {@link ReplicasVerify}. + */ +class TestReplicasVerify { + + @Test + void testRetriesWithForcedContainerCacheRefreshOnBlockNotFound() throws Exception { + ReplicasVerify replicasVerify = new ReplicasVerify(); + ReplicaVerifier blockExistenceVerifier = mock(ReplicaVerifier.class); + ReplicaVerifier checksumVerifier = mock(ReplicaVerifier.class); + + when(blockExistenceVerifier.getType()).thenReturn("blockExistence"); + when(checksumVerifier.getType()).thenReturn("checksum"); + when(blockExistenceVerifier.verifyBlock(any(), any())) + .thenReturn(BlockVerificationResult.failCheckAndRefreshKeyLocation("Block not found")) + .thenReturn(BlockVerificationResult.pass()); + when(checksumVerifier.verifyBlock(any(), any())) + .thenReturn(BlockVerificationResult.pass()) + .thenReturn(BlockVerificationResult.pass()); + setField(replicasVerify, "replicaVerifiers", Arrays.asList(blockExistenceVerifier, checksumVerifier)); + setField(replicasVerify, "replication", emptyReplicationFilter()); + + OzoneClient ozoneClient = mock(OzoneClient.class); + ClientProtocol proxy = mock(ClientProtocol.class); + when(ozoneClient.getProxy()).thenReturn(proxy); + + OmKeyLocationInfo firstLocation = createKeyLocationInfo(1L, 11L); + OmKeyLocationInfo refreshedLocation = createKeyLocationInfo(2L, 22L); + when(proxy.getKeyInfo("vol1", "bucket1", "key1", false)) + .thenReturn(createKeyInfo(firstLocation)); + when(proxy.getKeyInfo("vol1", "bucket1", "key1", true)) + .thenReturn(createKeyInfo(refreshedLocation)); + + ObjectNode root = JsonUtils.createObjectNode(null); + ArrayNode keysArray = root.putArray("keys"); + AtomicBoolean allKeysPassed = new AtomicBoolean(true); + + replicasVerify.processKey(ozoneClient, "vol1", "bucket1", "key1", keysArray, allKeysPassed); + + verify(proxy).getKeyInfo("vol1", "bucket1", "key1", false); + verify(proxy).getKeyInfo("vol1", "bucket1", "key1", true); + verify(blockExistenceVerifier, times(2)).verifyBlock(any(), any()); + verify(checksumVerifier, times(2)).verifyBlock(any(), any()); + assertTrue(allKeysPassed.get()); + assertEquals(0, keysArray.size()); + } + + @Test + void testDoesNotRefreshOnNonRefreshableFailure() throws Exception { + ReplicasVerify replicasVerify = new ReplicasVerify(); + ReplicaVerifier blockExistenceVerifier = mock(ReplicaVerifier.class); + ReplicaVerifier checksumVerifier = mock(ReplicaVerifier.class); + + when(blockExistenceVerifier.getType()).thenReturn("blockExistence"); + when(checksumVerifier.getType()).thenReturn("checksum"); + when(blockExistenceVerifier.verifyBlock(any(), any())) + .thenReturn(BlockVerificationResult.failCheck("Block does not exist on this replica")); + when(checksumVerifier.verifyBlock(any(), any())) + .thenReturn(BlockVerificationResult.pass()); + setField(replicasVerify, "replicaVerifiers", Arrays.asList(blockExistenceVerifier, checksumVerifier)); + setField(replicasVerify, "replication", emptyReplicationFilter()); + + OzoneClient ozoneClient = mock(OzoneClient.class); + ClientProtocol proxy = mock(ClientProtocol.class); + when(ozoneClient.getProxy()).thenReturn(proxy); + when(proxy.getKeyInfo("vol1", "bucket1", "key1", false)) + .thenReturn(createKeyInfo(createKeyLocationInfo(1L, 11L))); + + ObjectNode root = JsonUtils.createObjectNode(null); + ArrayNode keysArray = root.putArray("keys"); + AtomicBoolean allKeysPassed = new AtomicBoolean(true); + + replicasVerify.processKey(ozoneClient, "vol1", "bucket1", "key1", keysArray, allKeysPassed); + + verify(proxy).getKeyInfo("vol1", "bucket1", "key1", false); + verify(proxy, never()).getKeyInfo("vol1", "bucket1", "key1", true); + verify(blockExistenceVerifier, times(1)).verifyBlock(any(), any()); + verify(checksumVerifier, times(1)).verifyBlock(any(), any()); + assertFalse(allKeysPassed.get()); + assertEquals(1, keysArray.size()); + assertFalse(keysArray.get(0).get("pass").asBoolean()); + } + + private static OmKeyInfo createKeyInfo(OmKeyLocationInfo keyLocationInfo) { + OmKeyLocationInfoGroup latestVersionLocations = + new OmKeyLocationInfoGroup(0, Collections.singletonList(keyLocationInfo)); + return new OmKeyInfo.Builder() + .setVolumeName("vol1") + .setBucketName("bucket1") + .setKeyName("key1") + .setReplicationConfig(StandaloneReplicationConfig.getInstance(HddsProtos.ReplicationFactor.ONE)) + .setOmKeyLocationInfos(Collections.singletonList(latestVersionLocations)) + .build(); + } + + private static OmKeyLocationInfo createKeyLocationInfo(long containerId, long localId) { + DatanodeDetails datanode = MockDatanodeDetails.randomDatanodeDetails(); + Pipeline pipeline = Pipeline.newBuilder() + .setId(PipelineID.randomId()) + .setReplicationConfig(StandaloneReplicationConfig.getInstance(HddsProtos.ReplicationFactor.ONE)) + .setState(Pipeline.PipelineState.OPEN) + .setNodes(Collections.singletonList(datanode)) + .build(); + + return new OmKeyLocationInfo.Builder() + .setBlockID(new BlockID(containerId, localId)) + .setPipeline(pipeline) + .setLength(1) + .setOffset(0) + .setCreateVersion(0) + .build(); + } + + private static ShellReplicationOptions emptyReplicationFilter() { + ShellReplicationOptions replication = mock(ShellReplicationOptions.class); + when(replication.fromParams(any())).thenReturn(Optional.empty()); + return replication; + } + + private static void setField(Object target, String fieldName, Object value) throws Exception { + Class type = target.getClass(); + while (type != null) { + try { + Field field = type.getDeclaredField(fieldName); + field.setAccessible(true); + field.set(target, value); + return; + } catch (NoSuchFieldException ignored) { + type = type.getSuperclass(); + } + } + throw new IOException("Unable to set field '" + fieldName + "' on " + target.getClass().getName()); + } +} From 4549a3bf7988ab8a5971b125c35b4a2b7c62d0b7 Mon Sep 17 00:00:00 2001 From: Aryan Gupta Date: Tue, 18 Aug 2026 17:09:55 +0530 Subject: [PATCH 2/5] Fixed checkstyle. --- .../apache/hadoop/ozone/debug/replicas/ChecksumVerifier.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ChecksumVerifier.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ChecksumVerifier.java index c2fb110721f1..2acc514481a4 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ChecksumVerifier.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ChecksumVerifier.java @@ -23,6 +23,7 @@ import org.apache.commons.io.output.NullOutputStream; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.scm.OzoneClientConfig; import org.apache.hadoop.hdds.scm.XceiverClientManager; import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; @@ -71,7 +72,7 @@ public BlockVerificationResult verifyBlock(DatanodeDetails datanode, OmKeyLocati Throwable cause = e.getCause() != null ? e.getCause() : e; StorageContainerException storageContainerException = findStorageContainerException(e); if (storageContainerException != null - && storageContainerException.getResult() == org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.NO_SUCH_BLOCK) { + && storageContainerException.getResult() == ContainerProtos.Result.NO_SUCH_BLOCK) { return BlockVerificationResult.failCheckAndRefreshKeyLocation(storageContainerException.getMessage()); } if (cause instanceof OzoneChecksumException) { From ef685deeb1f014855b5eeb201573ee2eb1dc400e Mon Sep 17 00:00:00 2001 From: Aryan Gupta Date: Sat, 22 Aug 2026 12:31:42 +0530 Subject: [PATCH 3/5] Handle OM refresh failures in replicas verify retry path. --- .../ozone/debug/replicas/ReplicasVerify.java | 28 +++++++++++--- .../debug/replicas/TestReplicasVerify.java | 37 +++++++++++++++++++ 2 files changed, 60 insertions(+), 5 deletions(-) diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java index f7c3675434ee..ba43ca92a7ce 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java @@ -348,9 +348,15 @@ void processKey(OzoneClient ozoneClient, String volumeName, String bucketName, S KeyVerificationResult keyVerificationResult = verifyKey(keyInfo, volumeName, bucketName, keyName); if (keyVerificationResult.shouldRefreshKeyLocation()) { - OmKeyInfo refreshedKeyInfo = ozoneClient.getProxy().getKeyInfo( - volumeName, bucketName, keyName, true); - keyVerificationResult = verifyKey(refreshedKeyInfo, volumeName, bucketName, keyName); + try { + OmKeyInfo refreshedKeyInfo = ozoneClient.getProxy().getKeyInfo( + volumeName, bucketName, keyName, true); + keyVerificationResult = verifyKey(refreshedKeyInfo, volumeName, bucketName, keyName); + } catch (IOException e) { + LOG.warn("Unable to refresh key location from OM for key {}/{}/{}; marking verification incomplete.", + volumeName, bucketName, keyName, e); + keyVerificationResult.markRefreshFailure(e.getMessage()); + } } keyVerificationResult.getKeyNode().put("pass", keyVerificationResult.passed()); @@ -381,6 +387,7 @@ private KeyVerificationResult verifyKey( boolean keyPass = true; boolean shouldRefreshKeyLocation = false; Set failedVerificationTypes = new HashSet<>(); + List refreshChecks = new ArrayList<>(); for (OmKeyLocationInfo keyLocation : keyInfo.getLatestVersionLocations().getBlocksLatestVersionOnly()) { long containerID = keyLocation.getContainerID(); @@ -412,6 +419,7 @@ private KeyVerificationResult verifyKey( checkNode.put("pass", result.passed()); if (result.shouldRefreshKeyLocation()) { shouldRefreshKeyLocation = true; + refreshChecks.add(checkNode); } ArrayNode failuresArray = checkNode.putArray("failures"); @@ -436,7 +444,7 @@ private KeyVerificationResult verifyKey( } } - return new KeyVerificationResult(keyNode, keyPass, failedVerificationTypes, shouldRefreshKeyLocation); + return new KeyVerificationResult(keyNode, keyPass, failedVerificationTypes, shouldRefreshKeyLocation, refreshChecks); } private static final class KeyVerificationResult { @@ -444,13 +452,15 @@ private static final class KeyVerificationResult { private final boolean pass; private final Set failedVerificationTypes; private final boolean refreshKeyLocation; + private final List refreshChecks; private KeyVerificationResult(ObjectNode keyNode, boolean pass, Set failedVerificationTypes, - boolean refreshKeyLocation) { + boolean refreshKeyLocation, List refreshChecks) { this.keyNode = keyNode; this.pass = pass; this.failedVerificationTypes = failedVerificationTypes; this.refreshKeyLocation = refreshKeyLocation; + this.refreshChecks = refreshChecks; } private ObjectNode getKeyNode() { @@ -468,6 +478,14 @@ private Set getFailedVerificationTypes() { private boolean shouldRefreshKeyLocation() { return refreshKeyLocation; } + + private void markRefreshFailure(String message) { + String refreshFailure = "Failed to refresh key location from OM: " + message; + for (ObjectNode checkNode : refreshChecks) { + checkNode.put("completed", false); + checkNode.withArray("failures").addObject().put("message", refreshFailure); + } + } } /** diff --git a/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java index 9337393c24fb..a8d008140ecc 100644 --- a/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java +++ b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java @@ -154,6 +154,43 @@ void testDoesNotRefreshOnNonRefreshableFailure() throws Exception { assertFalse(keysArray.get(0).get("pass").asBoolean()); } + @Test + void testRefreshFailureMarksKeyIncompleteAndDoesNotAbortRun() throws Exception { + ReplicasVerify replicasVerify = new ReplicasVerify(); + ReplicaVerifier blockExistenceVerifier = mock(ReplicaVerifier.class); + + when(blockExistenceVerifier.getType()).thenReturn("blockExistence"); + when(blockExistenceVerifier.verifyBlock(any(), any())) + .thenReturn(BlockVerificationResult.failCheckAndRefreshKeyLocation("Block not found")); + setField(replicasVerify, "replicaVerifiers", Collections.singletonList(blockExistenceVerifier)); + setField(replicasVerify, "replication", emptyReplicationFilter()); + + OzoneClient ozoneClient = mock(OzoneClient.class); + ClientProtocol proxy = mock(ClientProtocol.class); + when(ozoneClient.getProxy()).thenReturn(proxy); + when(proxy.getKeyInfo("vol1", "bucket1", "key1", false)) + .thenReturn(createKeyInfo(createKeyLocationInfo(1L, 11L))); + when(proxy.getKeyInfo("vol1", "bucket1", "key1", true)) + .thenThrow(new IOException("OM refresh failed")); + + ObjectNode root = JsonUtils.createObjectNode(null); + ArrayNode keysArray = root.putArray("keys"); + AtomicBoolean allKeysPassed = new AtomicBoolean(true); + + replicasVerify.processKey(ozoneClient, "vol1", "bucket1", "key1", keysArray, allKeysPassed); + + verify(proxy).getKeyInfo("vol1", "bucket1", "key1", false); + verify(proxy).getKeyInfo("vol1", "bucket1", "key1", true); + assertFalse(allKeysPassed.get()); + assertEquals(1, keysArray.size()); + assertFalse(keysArray.get(0).get("pass").asBoolean()); + assertFalse(keysArray.get(0).get("blocks").get(0) + .get("replicas").get(0).get("checks").get(0).get("completed").asBoolean()); + assertThat(keysArray.get(0).get("blocks").get(0) + .get("replicas").get(0).get("checks").get(0).get("failures").toString()) + .contains("Failed to refresh key location from OM: OM refresh failed"); + } + private static boolean isRefreshContainerLocationsFromScmEnabled(ReplicasVerify command) throws Exception { Field field = ReplicasVerify.class.getDeclaredField("refreshContainerLocationsFromScm"); field.setAccessible(true); From 29a4d9ed3427d11bc04f5ea75f86f92f5a25aeb8 Mon Sep 17 00:00:00 2001 From: Aryan Gupta Date: Sat, 22 Aug 2026 12:37:33 +0530 Subject: [PATCH 4/5] Fixed Checkstyle. --- .../org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java index ba43ca92a7ce..a170eee37f22 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java @@ -444,7 +444,8 @@ private KeyVerificationResult verifyKey( } } - return new KeyVerificationResult(keyNode, keyPass, failedVerificationTypes, shouldRefreshKeyLocation, refreshChecks); + return new KeyVerificationResult( + keyNode, keyPass, failedVerificationTypes, shouldRefreshKeyLocation, refreshChecks); } private static final class KeyVerificationResult { From a50a2703f842fa4d586f107f046b27c2177a9b2f Mon Sep 17 00:00:00 2001 From: Aryan Gupta Date: Sat, 22 Aug 2026 14:54:58 +0530 Subject: [PATCH 5/5] Addressed comment. --- .../ozone/debug/replicas/ReplicasVerify.java | 27 +++++++++++++++++-- .../debug/replicas/TestReplicasVerify.java | 25 +++++++++++++++++ 2 files changed, 50 insertions(+), 2 deletions(-) diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java index a170eee37f22..784e3fdc9186 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/debug/replicas/ReplicasVerify.java @@ -338,8 +338,16 @@ void checkBucket(OzoneClient ozoneClient, OzoneBucket bucket, ArrayNode keysArra void processKey(OzoneClient ozoneClient, String volumeName, String bucketName, String keyName, ArrayNode keysArray, AtomicBoolean allKeysPassed) throws IOException { keysProcessed.incrementAndGet(); - OmKeyInfo keyInfo = ozoneClient.getProxy().getKeyInfo( - volumeName, bucketName, keyName, refreshContainerLocationsFromScm); + OmKeyInfo keyInfo; + try { + keyInfo = ozoneClient.getProxy().getKeyInfo( + volumeName, bucketName, keyName, refreshContainerLocationsFromScm); + } catch (IOException e) { + LOG.warn("Unable to fetch key info from OM for key {}/{}/{}; marking verification incomplete.", + volumeName, bucketName, keyName, e); + markKeyFetchFailure(volumeName, bucketName, keyName, e.getMessage(), keysArray, allKeysPassed); + return; + } // Check if key should be processed based on replication config if (!shouldProcessKeyByReplicationType(keyInfo)) { @@ -376,6 +384,21 @@ void processKey(OzoneClient ozoneClient, String volumeName, String bucketName, S } } + private void markKeyFetchFailure(String volumeName, String bucketName, String keyName, String message, + ArrayNode keysArray, AtomicBoolean allKeysPassed) { + ObjectNode keyNode = JsonUtils.createObjectNode(null); + keyNode.put("volumeName", volumeName); + keyNode.put("bucketName", bucketName); + keyNode.put("name", keyName); + keyNode.putArray("blocks"); + keyNode.put("completed", false); + keyNode.put("pass", false); + keyNode.putArray("failures").addObject().put("message", "Failed to fetch key info from OM: " + message); + keysFailed.incrementAndGet(); + allKeysPassed.set(false); + keysArray.add(keyNode); + } + private KeyVerificationResult verifyKey( OmKeyInfo keyInfo, String volumeName, String bucketName, String keyName) { ObjectNode keyNode = JsonUtils.createObjectNode(null); diff --git a/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java index a8d008140ecc..a70d7bc7bafc 100644 --- a/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java +++ b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/replicas/TestReplicasVerify.java @@ -191,6 +191,31 @@ void testRefreshFailureMarksKeyIncompleteAndDoesNotAbortRun() throws Exception { .contains("Failed to refresh key location from OM: OM refresh failed"); } + @Test + void testInitialGetKeyInfoFailureMarksKeyIncompleteAndDoesNotAbortRun() throws Exception { + ReplicasVerify replicasVerify = new ReplicasVerify(); + OzoneClient ozoneClient = mock(OzoneClient.class); + ClientProtocol proxy = mock(ClientProtocol.class); + when(ozoneClient.getProxy()).thenReturn(proxy); + when(proxy.getKeyInfo("vol1", "bucket1", "key1", false)) + .thenThrow(new IOException("Key not found")); + + ObjectNode root = JsonUtils.createObjectNode(null); + ArrayNode keysArray = root.putArray("keys"); + AtomicBoolean allKeysPassed = new AtomicBoolean(true); + + replicasVerify.processKey(ozoneClient, "vol1", "bucket1", "key1", keysArray, allKeysPassed); + + verify(proxy).getKeyInfo("vol1", "bucket1", "key1", false); + verify(proxy, never()).getKeyInfo("vol1", "bucket1", "key1", true); + assertFalse(allKeysPassed.get()); + assertEquals(1, keysArray.size()); + assertFalse(keysArray.get(0).get("completed").asBoolean()); + assertFalse(keysArray.get(0).get("pass").asBoolean()); + assertThat(keysArray.get(0).get("failures").toString()) + .contains("Failed to fetch key info from OM: Key not found"); + } + private static boolean isRefreshContainerLocationsFromScmEnabled(ReplicasVerify command) throws Exception { Field field = ReplicasVerify.class.getDeclaredField("refreshContainerLocationsFromScm"); field.setAccessible(true);