This is an automated email from the ASF dual-hosted git repository.
ArafatKhan2198 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new d9b6cd4fb31 HDDS-15863. Add Manual OM DB Rebuild Support for Recon
Bootstrapping. (#10825).
d9b6cd4fb31 is described below
commit d9b6cd4fb31c0a5328d8284e4063c2048d425c77
Author: Arafat2198 <[email protected]>
AuthorDate: Wed Jul 29 14:28:28 2026 +0530
HDDS-15863. Add Manual OM DB Rebuild Support for Recon Bootstrapping.
(#10825).
Co-authored-by: Cursor <[email protected]>
---
.../ozone/recon/TestReconWithOzoneManager.java | 116 +++++++++++++++++++++
.../ozone/recon/api/TriggerDBSyncEndpoint.java | 12 +++
.../types/OMDBReprocessResponse.java} | 55 +++++-----
.../recon/recovery/ReconOmMetadataManagerImpl.java | 1 +
.../recon/spi/OzoneManagerServiceProvider.java | 7 ++
.../spi/impl/OzoneManagerServiceProviderImpl.java | 22 ++++
.../ozone/recon/tasks/ReconTaskControllerImpl.java | 19 +++-
.../tasks/ReconTaskReInitializationEvent.java | 3 +-
.../impl/TestOzoneManagerServiceProviderImpl.java | 40 +++++++
.../recon/tasks/TestReconTaskControllerImpl.java | 22 +++-
10 files changed, 260 insertions(+), 37 deletions(-)
diff --git
a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconWithOzoneManager.java
b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconWithOzoneManager.java
index dfdc3de2412..8896f74c415 100644
---
a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconWithOzoneManager.java
+++
b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconWithOzoneManager.java
@@ -41,6 +41,7 @@
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicLong;
import org.apache.hadoop.hdds.JsonTestUtils;
import org.apache.hadoop.hdds.client.BlockID;
import org.apache.hadoop.hdds.client.StandaloneReplicationConfig;
@@ -52,9 +53,11 @@
import org.apache.hadoop.ozone.MiniOzoneCluster;
import org.apache.hadoop.ozone.om.OMMetadataManager;
import org.apache.hadoop.ozone.om.helpers.BucketLayout;
+import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
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.om.helpers.OmVolumeArgs;
import org.apache.hadoop.ozone.recon.metrics.OzoneManagerSyncMetrics;
import org.apache.hadoop.ozone.recon.spi.impl.OzoneManagerServiceProviderImpl;
import org.apache.http.HttpEntity;
@@ -74,6 +77,8 @@
* Test Ozone Recon.
*/
public class TestReconWithOzoneManager {
+ private static final AtomicLong OBJECT_ID_SEQUENCE = new AtomicLong();
+
private static MiniOzoneCluster cluster = null;
private static OMMetadataManager metadataManager;
private static CloseableHttpClient httpClient;
@@ -393,6 +398,93 @@ private long getReconTaskAttributeFromJson(String
taskStatusResponse,
* Helper function to add voli/bucketi/keyi to containeri to OM Metadata.
* For test purpose each container will have only one key.
*/
+ @Test
+ public void testManualOMDBRebuild() throws Exception {
+ // 1. Stop Recon
+ recon.stop();
+
+ // 2. Write keys to OM
+ addKeys(20, 25);
+ long omLatestSeqNumber = ((RDBStore) metadataManager.getStore())
+ .getDb().getLatestSequenceNumber();
+ java.io.File omDbDir = metadataManager.getStore().getDbLocation();
+
+ // 3. Stop OM (flush to disk)
+ cluster.getOzoneManager().stop();
+
+ // 4. Copy OM DB into Recon's OM snapshot dir. The dir and the
"om.snapshot.db_" prefix must
+ // match what ReconOmMetadataManagerImpl#start looks up via
ReconUtils#getLastKnownDB, otherwise
+ // Recon will not load the copied DB on restart.
+ java.io.File reconOmSnapshotDir = new ReconUtils()
+ .getReconDbDir(cluster.getConf(),
ReconServerConfigKeys.OZONE_RECON_OM_SNAPSHOT_DB_DIR);
+ java.io.File reconOmDbDir = new java.io.File(reconOmSnapshotDir,
+ ReconConstants.RECON_OM_SNAPSHOT_DB + "_" +
System.currentTimeMillis());
+ org.apache.commons.io.FileUtils.deleteDirectory(reconOmDbDir);
+ org.apache.commons.io.FileUtils.copyDirectory(omDbDir, reconOmDbDir);
+
+ // 5. Restart Recon and confirm it loaded the OM DB copy we placed above.
Shutting OM down
+ // writes its own trailing records, so the copy is at or ahead of the
sequence number we
+ // captured while OM was up.
+ recon.start(cluster.getConf());
+ OzoneManagerServiceProviderImpl reconImpl =
(OzoneManagerServiceProviderImpl)
+ recon.getReconServer().getOzoneManagerServiceProvider();
+ long reconLoadedSeqNumber = ((RDBStore)
reconImpl.getOMMetadataManagerInstance().getStore())
+ .getDb().getLatestSequenceNumber();
+ assertThat(reconLoadedSeqNumber).isGreaterThanOrEqualTo(omLatestSeqNumber);
+
+ // 6. POST reinit
+ String triggerUrl = "http://" +
cluster.getConf().get(OZONE_RECON_HTTP_ADDRESS_KEY) +
+ "/api/v1/triggerdbsync/om/reinit";
+ org.apache.http.client.methods.HttpPost httpPost = new
org.apache.http.client.methods.HttpPost(triggerUrl);
+ HttpResponse response = httpClient.execute(httpPost);
+ assertEquals(202, response.getStatusLine().getStatusCode());
+
+ // 7. Wait for reprocess to succeed (REPROCESS_STAGING seq is set by
reInitializeTasks)
+ GenericTestUtils.waitFor(() -> {
+ try {
+ String taskStatusResponse = makeHttpCall(taskStatusURL);
+ long reconLatestSeqNumber = getReconTaskAttributeFromJson(
+ taskStatusResponse,
+ "REPROCESS_STAGING",
+ "lastUpdatedSeqNumber");
+ return reconLatestSeqNumber == reconLoadedSeqNumber;
+ } catch (Exception e) {
+ return false;
+ }
+ }, 1000, 30000);
+
+ // 8. Start OM
+ cluster.getOzoneManager().restart();
+ // restart() rebuilds the OM metadata manager, so re-fetch the live handle.
+ refreshOmMetadataManager();
+
+ // 9. Write more keys and verify delta resumes
+ addKeys(25, 30);
+ long newOmLatestSeqNumber = ((RDBStore)
cluster.getOzoneManager().getMetadataManager().getStore())
+ .getDb().getLatestSequenceNumber();
+
+ OzoneManagerServiceProviderImpl newImpl = (OzoneManagerServiceProviderImpl)
+ recon.getReconServer().getOzoneManagerServiceProvider();
+ newImpl.syncDataFromOM();
+
+ GenericTestUtils.waitFor(() -> {
+ try {
+ String taskStatusResponse = makeHttpCall(taskStatusURL);
+ long reconLatestSeqNumber = getReconTaskAttributeFromJson(
+ taskStatusResponse,
+ OmSnapshotRequest.name(),
+ "lastUpdatedSeqNumber");
+ return reconLatestSeqNumber == newOmLatestSeqNumber;
+ } catch (Exception e) {
+ return false;
+ }
+ }, 1000, 30000);
+ }
+
+ private static void refreshOmMetadataManager() {
+ metadataManager = cluster.getOzoneManager().getMetadataManager();
+ }
+
private void addKeys(int start, int end) throws Exception {
for (int i = start; i < end; i++) {
Pipeline pipeline = HddsTestUtils.getRandomPipeline();
@@ -421,6 +513,30 @@ private static void writeDataToOm(String key, String
bucket, String volume,
omKeyLocationInfoGroupList)
throws IOException {
+ // Recon's full reprocess resolves the parent bucket of every key it
reads, so the volume
+ // and bucket rows have to exist next to the key entry.
+ String volumeKey = metadataManager.getVolumeKey(volume);
+ if (metadataManager.getVolumeTable().get(volumeKey) == null) {
+ metadataManager.getVolumeTable().put(volumeKey,
+ OmVolumeArgs.newBuilder()
+ .setVolume(volume)
+ .setAdminName("TestUser")
+ .setOwnerName("TestUser")
+ .setObjectID(OBJECT_ID_SEQUENCE.incrementAndGet())
+ .build());
+ }
+
+ String bucketKey = metadataManager.getBucketKey(volume, bucket);
+ if (metadataManager.getBucketTable().get(bucketKey) == null) {
+ metadataManager.getBucketTable().put(bucketKey,
+ OmBucketInfo.newBuilder()
+ .setVolumeName(volume)
+ .setBucketName(bucket)
+ .setBucketLayout(getBucketLayout())
+ .setObjectID(OBJECT_ID_SEQUENCE.incrementAndGet())
+ .build());
+ }
+
String omKey = metadataManager.getOzoneKey(volume,
bucket, key);
diff --git
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/TriggerDBSyncEndpoint.java
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/TriggerDBSyncEndpoint.java
index 93e0f0c9231..493f2bfe4cd 100644
---
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/TriggerDBSyncEndpoint.java
+++
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/TriggerDBSyncEndpoint.java
@@ -25,6 +25,7 @@
import javax.ws.rs.core.MediaType;
import javax.ws.rs.core.Response;
import org.apache.hadoop.hdds.scm.server.OzoneStorageContainerManager;
+import org.apache.hadoop.ozone.recon.api.types.OMDBReprocessResponse;
import org.apache.hadoop.ozone.recon.scm.ReconStorageContainerManagerFacade;
import org.apache.hadoop.ozone.recon.spi.OzoneManagerServiceProvider;
@@ -54,6 +55,17 @@ public Response triggerOMDBSync() {
return Response.ok(isSuccess).build();
}
+ @POST
+ @Path("om/reinit")
+ public Response triggerOMDBReinit() {
+ OMDBReprocessResponse response =
ozoneManagerServiceProvider.triggerTaskRebuild();
+ if (response.getStatus() == OMDBReprocessResponse.Status.ACCEPTED) {
+ return Response.accepted(response).build();
+ } else {
+ return
Response.status(Response.Status.CONFLICT).entity(response).build();
+ }
+ }
+
@POST
@Path("scm/snapshot")
public Response triggerSCMDBSnapshotSync() {
diff --git
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/OzoneManagerServiceProvider.java
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/types/OMDBReprocessResponse.java
similarity index 56%
copy from
hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/OzoneManagerServiceProvider.java
copy to
hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/types/OMDBReprocessResponse.java
index dc32b1692dd..3988803d764 100644
---
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/OzoneManagerServiceProvider.java
+++
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/types/OMDBReprocessResponse.java
@@ -15,34 +15,37 @@
* limitations under the License.
*/
-package org.apache.hadoop.ozone.recon.spi;
+package org.apache.hadoop.ozone.recon.api.types;
-import org.apache.hadoop.ozone.om.OMMetadataManager;
+import com.fasterxml.jackson.annotation.JsonProperty;
/**
- * Interface to access OM endpoints.
+ * Response for OM DB manual reprocess request.
*/
-public interface OzoneManagerServiceProvider {
-
- /**
- * Start a task to sync data from OM.
- */
- void start();
-
- /**
- * Stop the OM sync data.
- */
- void stop() throws Exception;
-
- /**
- * Return instance of OM Metadata manager.
- * @return OM metadata manager instance.
- */
- OMMetadataManager getOMMetadataManagerInstance();
-
- /**
- * Trigger a sync between Recon and OM.
- * @return whether the trigger happened or not
- */
- boolean triggerSyncDataFromOMImmediately();
+public class OMDBReprocessResponse {
+
+ @JsonProperty("status")
+ private Status status;
+
+ @JsonProperty("message")
+ private String message;
+
+ /** Result of a manual OM DB reprocess request. */
+ public enum Status {
+ ACCEPTED,
+ RETRY
+ }
+
+ public OMDBReprocessResponse(Status status, String message) {
+ this.status = status;
+ this.message = message;
+ }
+
+ public Status getStatus() {
+ return status;
+ }
+
+ public String getMessage() {
+ return message;
+ }
}
diff --git
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/recovery/ReconOmMetadataManagerImpl.java
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/recovery/ReconOmMetadataManagerImpl.java
index ae0059f039d..d0f92d0485f 100644
---
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/recovery/ReconOmMetadataManagerImpl.java
+++
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/recovery/ReconOmMetadataManagerImpl.java
@@ -104,6 +104,7 @@ public void start(OzoneConfiguration configuration) throws
IOException {
LOG.info("Starting ReconOMMetadataManagerImpl");
File reconDbDir =
reconUtils.getReconDbDir(configuration,
OZONE_RECON_OM_SNAPSHOT_DB_DIR);
+ LOG.info("reconDbDir is: {}", reconDbDir.getAbsolutePath());
File lastKnownOMSnapshot =
reconUtils.getLastKnownDB(reconDbDir, RECON_OM_SNAPSHOT_DB);
if (lastKnownOMSnapshot != null) {
diff --git
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/OzoneManagerServiceProvider.java
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/OzoneManagerServiceProvider.java
index dc32b1692dd..1bc45c6627a 100644
---
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/OzoneManagerServiceProvider.java
+++
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/OzoneManagerServiceProvider.java
@@ -18,6 +18,7 @@
package org.apache.hadoop.ozone.recon.spi;
import org.apache.hadoop.ozone.om.OMMetadataManager;
+import org.apache.hadoop.ozone.recon.api.types.OMDBReprocessResponse;
/**
* Interface to access OM endpoints.
@@ -45,4 +46,10 @@ public interface OzoneManagerServiceProvider {
* @return whether the trigger happened or not
*/
boolean triggerSyncDataFromOMImmediately();
+
+ /**
+ * Trigger the OM DB rebuild process.
+ * @return OMDBReprocessResponse containing the status of the request.
+ */
+ OMDBReprocessResponse triggerTaskRebuild();
}
diff --git
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/impl/OzoneManagerServiceProviderImpl.java
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/impl/OzoneManagerServiceProviderImpl.java
index bec49088d27..d221168aa00 100644
---
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/impl/OzoneManagerServiceProviderImpl.java
+++
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/impl/OzoneManagerServiceProviderImpl.java
@@ -85,6 +85,7 @@
import org.apache.hadoop.ozone.recon.ReconContext;
import org.apache.hadoop.ozone.recon.ReconServerConfigKeys;
import org.apache.hadoop.ozone.recon.ReconUtils;
+import org.apache.hadoop.ozone.recon.api.types.OMDBReprocessResponse;
import org.apache.hadoop.ozone.recon.metrics.OzoneManagerSyncMetrics;
import org.apache.hadoop.ozone.recon.metrics.ReconSyncMetrics;
import org.apache.hadoop.ozone.recon.recovery.ReconOMMetadataManager;
@@ -562,6 +563,27 @@ boolean updateReconOmDBWithNewSnapshot() throws
IOException {
}
}
+ @Override
+ public OMDBReprocessResponse triggerTaskRebuild() {
+ if (omMetadataManager == null || omMetadataManager.getStore() == null) {
+ return new OMDBReprocessResponse(OMDBReprocessResponse.Status.RETRY,
+ "Recon has not loaded an OM DB yet, so there is nothing to rebuild.
Ensure an OM DB snapshot is "
+ + "present in the Recon OM DB directory and has been loaded,
then retry.");
+ }
+
+ ReconTaskController.ReInitializationResult result =
reconTaskController.queueReInitializationEvent(
+
ReconTaskReInitializationEvent.ReInitializationReason.MANUAL_OM_DB_REBUILD);
+
+ if (result == ReconTaskController.ReInitializationResult.SUCCESS) {
+ return new OMDBReprocessResponse(OMDBReprocessResponse.Status.ACCEPTED,
+ "Manual OM DB rebuild queued successfully.");
+ } else {
+ return new OMDBReprocessResponse(OMDBReprocessResponse.Status.RETRY,
+ "Manual OM DB rebuild could not be queued. Buffer might be full or
another rebuild is "
+ + "pending. Please retry.");
+ }
+ }
+
/**
* Get Delta updates from OM through RPC call and apply to local OM DB as
* well as accumulate in a buffer.
diff --git
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/ReconTaskControllerImpl.java
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/ReconTaskControllerImpl.java
index 223d7583d04..5b142c0189d 100644
---
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/ReconTaskControllerImpl.java
+++
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/ReconTaskControllerImpl.java
@@ -617,6 +617,10 @@ public synchronized
ReconTaskController.ReInitializationResult queueReInitializa
// Track reprocess submission
controllerMetrics.incrTotalReprocessSubmittedToQueue();
+ if (reason ==
ReconTaskReInitializationEvent.ReInitializationReason.MANUAL_OM_DB_REBUILD) {
+ lastRetryTimestamp.set(0);
+ }
+
ReInitializationResult reInitializationResult =
validateRetryCountAndDelay();
if (null != reInitializationResult) {
return reInitializationResult;
@@ -746,9 +750,14 @@ public ReconOMMetadataManager
createOMCheckpoint(ReconOMMetadataManager omMetaMa
// Create temporary directory for checkpoint
String parentPath = cleanTempCheckPointPath(omMetaManager);
- // Create checkpoint
+ // Create checkpoint. getCheckpoint returns null when RocksDB fails to
snapshot
+ // (e.g. a manually placed OM DB that is incomplete or corrupt).
DBCheckpoint checkpoint =
omMetaManager.getStore().getCheckpoint(parentPath, true);
-
+ if (checkpoint == null) {
+ throw new IOException("Failed to create OM DB checkpoint at " +
parentPath
+ + "; the on-disk OM DB may be incomplete or corrupt.");
+ }
+
return omMetaManager.createCheckpointReconMetadataManager(configuration,
checkpoint);
}
@@ -840,7 +849,7 @@ public long getDroppedBatches() {
* Reset retry counters - for testing purposes.
*/
@VisibleForTesting
- void resetRetryCounters() {
+ public void resetRetryCounters() {
eventProcessRetryCount.set(0);
lastRetryTimestamp.set(0);
}
@@ -868,8 +877,8 @@ AtomicBoolean getTasksFailedFlag() {
*/
private void cleanupPreExistingCheckpoints() {
try {
- if (currentOMMetadataManager == null) {
- LOG.debug("No current OM metadata manager, skipping pre-existing
checkpoint cleanup");
+ if (currentOMMetadataManager == null ||
currentOMMetadataManager.getStore() == null) {
+ LOG.debug("No current OM metadata manager or store, skipping
pre-existing checkpoint cleanup");
return;
}
diff --git
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/ReconTaskReInitializationEvent.java
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/ReconTaskReInitializationEvent.java
index e241c48d11b..be895129a00 100644
---
a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/ReconTaskReInitializationEvent.java
+++
b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/tasks/ReconTaskReInitializationEvent.java
@@ -36,7 +36,8 @@ public class ReconTaskReInitializationEvent implements
ReconEvent {
public enum ReInitializationReason {
BUFFER_OVERFLOW,
TASK_FAILURES,
- MANUAL_TRIGGER
+ MANUAL_TRIGGER,
+ MANUAL_OM_DB_REBUILD
}
public ReconTaskReInitializationEvent(ReInitializationReason reason,
diff --git
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/spi/impl/TestOzoneManagerServiceProviderImpl.java
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/spi/impl/TestOzoneManagerServiceProviderImpl.java
index c5664934116..de28008c710 100644
---
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/spi/impl/TestOzoneManagerServiceProviderImpl.java
+++
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/spi/impl/TestOzoneManagerServiceProviderImpl.java
@@ -39,6 +39,7 @@
import static org.mockito.Mockito.doCallRealMethod;
import static org.mockito.Mockito.doNothing;
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;
@@ -64,6 +65,7 @@
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import org.apache.hadoop.ozone.recon.ReconContext;
import org.apache.hadoop.ozone.recon.ReconUtils;
+import org.apache.hadoop.ozone.recon.api.types.OMDBReprocessResponse;
import org.apache.hadoop.ozone.recon.common.ReconTestUtils;
import org.apache.hadoop.ozone.recon.metrics.OzoneManagerSyncMetrics;
import org.apache.hadoop.ozone.recon.recovery.ReconOMMetadataManager;
@@ -677,6 +679,44 @@ private ReconTaskStatusUpdaterManager
getMockTaskStatusUpdaterManager() {
return reconTaskStatusUpdaterManager;
}
+ @Test
+ public void testTriggerTaskRebuild() throws Exception {
+ ReconTaskController reconTaskController = mock(ReconTaskController.class);
+ ReconTaskStatusUpdaterManager taskStatusUpdaterManager =
mock(ReconTaskStatusUpdaterManager.class);
+ ReconUtils reconUtils = new ReconUtils();
+
+ // Test 1: Successful queue (OM DB is loaded)
+ ReconOMMetadataManager loadedOmMgr = mock(ReconOMMetadataManager.class);
+ when(loadedOmMgr.getStore()).thenReturn(mock(RDBStore.class));
+ OzoneManagerServiceProviderImpl serviceProvider = new
OzoneManagerServiceProviderImpl(
+ configuration, loadedOmMgr, reconTaskController,
+ reconUtils, ozoneManagerProtocol, reconContext,
taskStatusUpdaterManager);
+
+ when(reconTaskController.queueReInitializationEvent(
+
ReconTaskReInitializationEvent.ReInitializationReason.MANUAL_OM_DB_REBUILD))
+ .thenReturn(ReconTaskController.ReInitializationResult.SUCCESS);
+
+ OMDBReprocessResponse response = serviceProvider.triggerTaskRebuild();
+ assertEquals(OMDBReprocessResponse.Status.ACCEPTED, response.getStatus());
+
+ // Test 2: Buffer full / retry
+ when(reconTaskController.queueReInitializationEvent(
+
ReconTaskReInitializationEvent.ReInitializationReason.MANUAL_OM_DB_REBUILD))
+ .thenReturn(ReconTaskController.ReInitializationResult.RETRY_LATER);
+
+ response = serviceProvider.triggerTaskRebuild();
+ assertEquals(OMDBReprocessResponse.Status.RETRY, response.getStatus());
+
+ // Test 3: No OM DB loaded (store is null) -> RETRY, nothing queued
+ ReconTaskController noDbController = mock(ReconTaskController.class);
+ serviceProvider = new OzoneManagerServiceProviderImpl(
+ configuration, mock(ReconOMMetadataManager.class), noDbController,
+ reconUtils, ozoneManagerProtocol, reconContext,
taskStatusUpdaterManager);
+ response = serviceProvider.triggerTaskRebuild();
+ assertEquals(OMDBReprocessResponse.Status.RETRY, response.getStatus());
+ verify(noDbController, never()).queueReInitializationEvent(any());
+ }
+
private BucketLayout getBucketLayout() {
return BucketLayout.DEFAULT;
}
diff --git
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/tasks/TestReconTaskControllerImpl.java
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/tasks/TestReconTaskControllerImpl.java
index 5fefd4ee8d2..e9da0f76a0e 100644
---
a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/tasks/TestReconTaskControllerImpl.java
+++
b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/tasks/TestReconTaskControllerImpl.java
@@ -21,6 +21,7 @@
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyMap;
@@ -813,11 +814,22 @@ public void
testProcessReInitializationEventWithCheckpointedManager() throws Exc
verify(mockCheckpointedManager, times(1)).close();
}
- /**
- * Helper method for getting a mocked Task.
- * @param taskName name of the task.
- * @return instance of reconOmTask.
- */
+ @Test
+ public void testCreateOMCheckpointThrowsWhenCheckpointNull() throws
Exception {
+ // getStore().getCheckpoint() returns null when RocksDB fails to snapshot
an
+ // incomplete/corrupt on-disk OM DB; createOMCheckpoint must surface this
as an
+ // IOException so the caller handles it gracefully instead of NPE-ing.
+ ReconOMMetadataManager omMetadataManager =
mock(ReconOMMetadataManager.class);
+ DBStore dbStore = mock(DBStore.class);
+ when(omMetadataManager.getStore()).thenReturn(dbStore);
+ File tempDir = Paths.get(System.getProperty("java.io.tmpdir"),
"recon-test").toFile();
+ when(dbStore.getDbLocation()).thenReturn(tempDir);
+ when(dbStore.getCheckpoint(anyString(),
any(Boolean.class))).thenReturn(null);
+
+ ReconTaskControllerImpl controller = (ReconTaskControllerImpl)
reconTaskController;
+ assertThrows(IOException.class, () ->
controller.createOMCheckpoint(omMetadataManager));
+ }
+
private ReconOmTask getMockTask(String taskName) {
ReconOmTask reconOmTaskMock = mock(ReconOmTask.class);
when(reconOmTaskMock.getTaskName()).thenReturn(taskName);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]