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]

Reply via email to