This is an automated email from the ASF dual-hosted git repository.

jojochuang 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 d14585c8ce7 HDDS-15204. Quota repair includes snapshot pending-delete 
usage. (#10217)
d14585c8ce7 is described below

commit d14585c8ce73544416448f27fc756e3cfc565e69
Author: Wei-Chiu Chuang <[email protected]>
AuthorDate: Mon Jun 22 12:11:55 2026 -0700

    HDDS-15204. Quota repair includes snapshot pending-delete usage. (#10217)
    
    Generated-by: Cursor <[email protected]>
    Co-authored-by: sadanand48 <[email protected]>
---
 .../hadoop/ozone/om/helpers/OmBucketInfo.java      |   4 +-
 .../src/main/proto/OmClientProtocol.proto          |   2 +
 .../om/request/volume/OMQuotaRepairRequest.java    |   6 +
 .../hadoop/ozone/om/service/QuotaRepairTask.java   | 267 ++++++++++++++++++++-
 .../ozone/om/service/TestQuotaRepairTask.java      | 142 ++++++++++-
 5 files changed, 408 insertions(+), 13 deletions(-)

diff --git 
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmBucketInfo.java
 
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmBucketInfo.java
index 0d4eabf0bc2..8da2be2755a 100644
--- 
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmBucketInfo.java
+++ 
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmBucketInfo.java
@@ -267,7 +267,7 @@ public void decrUsedBytes(long bytes, boolean 
increasePendingDeleteBytes) {
     }
   }
 
-  private void incrSnapshotUsedBytes(long bytes) {
+  public void incrSnapshotUsedBytes(long bytes) {
     this.snapshotUsedBytes += bytes;
   }
 
@@ -282,7 +282,7 @@ public void decrUsedNamespace(long namespaceToUse, boolean 
increasePendingDelete
     }
   }
 
-  private void incrSnapshotUsedNamespace(long namespaceToUse) {
+  public void incrSnapshotUsedNamespace(long namespaceToUse) {
     this.snapshotUsedNamespace += namespaceToUse;
   }
 
diff --git 
a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto 
b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
index d71ff27f302..f50d812e51f 100644
--- a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
+++ b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
@@ -2381,6 +2381,8 @@ message BucketQuotaCount {
     required int64 diffUsedBytes = 3;
     required int64 diffUsedNamespace = 4;
     required bool supportOldQuota = 5 [default=false];
+    optional int64 diffSnapshotUsedBytes = 6;
+    optional int64 diffSnapshotUsedNamespace = 7;
 }
 
 message QuotaRepairResponse {
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/volume/OMQuotaRepairRequest.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/volume/OMQuotaRepairRequest.java
index 08b38cb2174..c58b700529f 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/volume/OMQuotaRepairRequest.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/volume/OMQuotaRepairRequest.java
@@ -124,6 +124,12 @@ private void updateBucketInfo(
       }
       bucketInfo.incrUsedBytes(bucketCountInfo.getDiffUsedBytes());
       bucketInfo.incrUsedNamespace(bucketCountInfo.getDiffUsedNamespace());
+      if (bucketCountInfo.hasDiffSnapshotUsedBytes()) {
+        
bucketInfo.incrSnapshotUsedBytes(bucketCountInfo.getDiffSnapshotUsedBytes());
+      }
+      if (bucketCountInfo.hasDiffSnapshotUsedNamespace()) {
+        
bucketInfo.incrSnapshotUsedNamespace(bucketCountInfo.getDiffSnapshotUsedNamespace());
+      }
       if (bucketCountInfo.getSupportOldQuota()) {
         OmBucketInfo.Builder builder = bucketInfo.toBuilder();
         if (bucketInfo.getQuotaInBytes() == OLD_QUOTA_DEFAULT) {
diff --git 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/QuotaRepairTask.java
 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/QuotaRepairTask.java
index cba6ad4aec2..2805f84c837 100644
--- 
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/QuotaRepairTask.java
+++ 
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/QuotaRepairTask.java
@@ -21,6 +21,7 @@
 import static org.apache.hadoop.hdds.utils.db.IteratorType.KEY_ONLY;
 import static org.apache.hadoop.ozone.OzoneConsts.OLD_QUOTA_DEFAULT;
 import static org.apache.hadoop.ozone.OzoneConsts.OM_KEY_PREFIX;
+import static 
org.apache.hadoop.ozone.om.helpers.SnapshotInfo.SnapshotStatus.SNAPSHOT_ACTIVE;
 
 import com.google.common.util.concurrent.UncheckedExecutionException;
 import com.google.protobuf.ServiceException;
@@ -31,8 +32,12 @@
 import java.util.ArrayList;
 import java.util.Collections;
 import java.util.HashMap;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
+import java.util.UUID;
 import java.util.concurrent.ArrayBlockingQueue;
 import java.util.concurrent.BlockingQueue;
 import java.util.concurrent.CompletableFuture;
@@ -52,15 +57,22 @@
 import org.apache.hadoop.hdds.utils.db.TableIterator;
 import org.apache.hadoop.ozone.om.OMMetadataManager;
 import org.apache.hadoop.ozone.om.OmMetadataManagerImpl;
+import org.apache.hadoop.ozone.om.OmSnapshot;
+import org.apache.hadoop.ozone.om.OmSnapshotManager;
 import org.apache.hadoop.ozone.om.OzoneManager;
+import org.apache.hadoop.ozone.om.SnapshotChainInfo;
+import org.apache.hadoop.ozone.om.SnapshotChainManager;
 import org.apache.hadoop.ozone.om.exceptions.OMException;
 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.RepeatedOmKeyInfo;
+import org.apache.hadoop.ozone.om.helpers.SnapshotInfo;
 import org.apache.hadoop.ozone.om.ratis.utils.OzoneManagerRatisUtils;
 import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
 import org.apache.hadoop.util.Time;
 import org.apache.ratis.protocol.ClientId;
+import org.apache.ratis.util.function.UncheckedAutoCloseableSupplier;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -72,6 +84,11 @@ public class QuotaRepairTask {
       QuotaRepairTask.class);
   private static final int BATCH_SIZE = 5000;
   private static final int TASK_THREAD_CNT = 3;
+  /**
+   * Parallel full-table scans: OBS keys, FSO files, dirs, active deleted 
keys/dirs,
+   * snapshot DB deleted keys/dirs.
+   */
+  private static final int QUOTA_REPAIR_SCAN_TASKS = 6;
   private static final AtomicBoolean IN_PROGRESS = new AtomicBoolean(false);
   private static final RepairStatus REPAIR_STATUS = new RepairStatus();
   private static final AtomicLong RUN_CNT = new AtomicLong(0);
@@ -102,8 +119,8 @@ public static String getStatus() {
 
   private boolean repairTask(List<String> buckets) {
     LOG.info("Starting quota repair task {}", REPAIR_STATUS);
-    // thread pool with 3 Table type * (1 task each + 3 thread for each task)
-    executor = Executors.newFixedThreadPool(3 * (1 + TASK_THREAD_CNT));
+    // thread pool: scan task types * (1 coordinator + worker threads per task)
+    executor = Executors.newFixedThreadPool(QUOTA_REPAIR_SCAN_TASKS * (1 + 
TASK_THREAD_CNT));
     try (OMMetadataManager activeMetaManager =
         createActiveDBCheckpoint(om.getMetadataManager(), 
om.getConfiguration())) {
       OzoneManagerProtocolProtos.QuotaRepairRequest.Builder builder
@@ -111,8 +128,6 @@ private boolean repairTask(List<String> buckets) {
       // repair active db
       repairActiveDb(activeMetaManager, builder, buckets);
 
-      // TODO: repair snapshots for quota
-
       // submit request to update
       ClientId clientId = ClientId.randomId();
       OzoneManagerProtocolProtos.OMRequest omRequest = 
OzoneManagerProtocolProtos.OMRequest.newBuilder()
@@ -175,6 +190,10 @@ private void repairActiveDb(
       bucketCountBuilder.setDiffUsedBytes(updatedBuckedInfo.getUsedBytes() - 
oriBucketInfo.getUsedBytes());
       bucketCountBuilder.setDiffUsedNamespace(
           updatedBuckedInfo.getUsedNamespace() - 
oriBucketInfo.getUsedNamespace());
+      bucketCountBuilder.setDiffSnapshotUsedBytes(
+          updatedBuckedInfo.getSnapshotUsedBytes() - 
oriBucketInfo.getSnapshotUsedBytes());
+      bucketCountBuilder.setDiffSnapshotUsedNamespace(
+          updatedBuckedInfo.getSnapshotUsedNamespace() - 
oriBucketInfo.getSnapshotUsedNamespace());
       bucketCountBuilder.setSupportOldQuota(oldQuota);
       builder.addBucketCount(bucketCountBuilder.build());
     }
@@ -253,17 +272,28 @@ private static void populateBucket(
     oriBucketInfoMap.put(bucketNameKey, bucketInfo.copyObject());
     bucketInfo.decrUsedBytes(bucketInfo.getUsedBytes(), false);
     bucketInfo.decrUsedNamespace(bucketInfo.getUsedNamespace(), false);
+    resetSnapshotBucketQuota(bucketInfo);
     nameBucketInfoMap.put(bucketNameKey, bucketInfo);
     
idBucketInfoMap.put(buildIdPath(metadataManager.getVolumeId(bucketInfo.getVolumeName()),
             bucketInfo.getObjectID()), bucketInfo);
   }
 
   private boolean isChange(OmBucketInfo lBucketInfo, OmBucketInfo rBucketInfo) 
{
-    if (lBucketInfo.getUsedNamespace() != rBucketInfo.getUsedNamespace()
-        || lBucketInfo.getUsedBytes() != rBucketInfo.getUsedBytes()) {
-      return true;
+    return lBucketInfo.getUsedNamespace() != rBucketInfo.getUsedNamespace()
+        || lBucketInfo.getUsedBytes() != rBucketInfo.getUsedBytes()
+        || lBucketInfo.getSnapshotUsedBytes() != 
rBucketInfo.getSnapshotUsedBytes()
+        || lBucketInfo.getSnapshotUsedNamespace() != 
rBucketInfo.getSnapshotUsedNamespace();
+  }
+
+  private static void resetSnapshotBucketQuota(OmBucketInfo bucketInfo) {
+    long snapBytes = bucketInfo.getSnapshotUsedBytes();
+    if (snapBytes != 0) {
+      bucketInfo.purgeSnapshotUsedBytes(snapBytes);
+    }
+    long snapNs = bucketInfo.getSnapshotUsedNamespace();
+    if (snapNs != 0) {
+      bucketInfo.purgeSnapshotUsedNamespace(snapNs);
     }
-    return false;
   }
   
   private static String buildNamePath(String volumeName, String bucketName) {
@@ -293,6 +323,8 @@ private void repairCount(
     Map<String, CountPair> keyCountMap = new ConcurrentHashMap<>();
     Map<String, CountPair> fileCountMap = new ConcurrentHashMap<>();
     Map<String, CountPair> directoryCountMap = new ConcurrentHashMap<>();
+    Map<String, CountPair> snapshotDeletedKeyMap = new ConcurrentHashMap<>();
+    Map<String, CountPair> snapshotDeletedDirMap = new ConcurrentHashMap<>();
     try {
       nameBucketInfoMap.keySet().stream().forEach(e -> keyCountMap.put(e,
           new CountPair()));
@@ -300,7 +332,9 @@ private void repairCount(
           new CountPair()));
       idBucketInfoMap.keySet().stream().forEach(e -> directoryCountMap.put(e,
           new CountPair()));
-      
+      nameBucketInfoMap.keySet().forEach(k -> snapshotDeletedKeyMap.put(k, new 
CountPair()));
+      idBucketInfoMap.keySet().forEach(k -> snapshotDeletedDirMap.put(k, new 
CountPair()));
+
       List<Future<?>> tasks = new ArrayList<>();
       tasks.add(executor.submit(() -> recalculateUsages(
           metadataManager.getKeyTable(BucketLayout.OBJECT_STORE),
@@ -311,6 +345,24 @@ private void repairCount(
       tasks.add(executor.submit(() -> recalculateUsages(
           metadataManager.getDirectoryTable(),
           directoryCountMap, "Directory usages", false)));
+
+      Map<Long, OmBucketInfo> bucketById = 
buildBucketByObjectId(idBucketInfoMap);
+
+      tasks.add(executor.submit(() -> recalculateDeletedKeyUsages(
+          metadataManager.getDeletedTable(), bucketById, snapshotDeletedKeyMap,
+          "active DB checkpoint")));
+      tasks.add(executor.submit(() -> recalculateDeletedDirNamespace(
+          metadataManager.getDeletedDirTable(), snapshotDeletedDirMap,
+          "active DB checkpoint")));
+      tasks.add(executor.submit(() -> {
+        try {
+          recalculateSnapshotDbPendingDeleteQuota(nameBucketInfoMap, 
bucketById, metadataManager,
+              snapshotDeletedKeyMap, snapshotDeletedDirMap);
+        } catch (IOException ex) {
+          throw new UncheckedIOException(ex);
+        }
+      }));
+
       for (Future<?> f : tasks) {
         f.get();
       }
@@ -326,9 +378,200 @@ private void repairCount(
     updateCountToBucketInfo(nameBucketInfoMap, keyCountMap);
     updateCountToBucketInfo(idBucketInfoMap, fileCountMap);
     updateCountToBucketInfo(idBucketInfoMap, directoryCountMap);
+    mergeSnapshotDeletedTableCounts(nameBucketInfoMap, snapshotDeletedKeyMap);
+    mergeDeletedDirSnapshotNamespace(idBucketInfoMap, snapshotDeletedDirMap);
     LOG.info("Completed quota repair counting for all keys, files and 
directories");
   }
 
+  private static Map<Long, OmBucketInfo> buildBucketByObjectId(
+      Map<String, OmBucketInfo> idBucketInfoMap) {
+    Map<Long, OmBucketInfo> bucketById = new HashMap<>();
+    for (OmBucketInfo bucketInfo : idBucketInfoMap.values()) {
+      bucketById.putIfAbsent(bucketInfo.getObjectID(), bucketInfo);
+    }
+    return bucketById;
+  }
+
+  /**
+   * Recompute pending-delete quota from each ACTIVE snapshot DB on repaired 
buckets' path chains.
+   * Runs as an executor task in parallel with active-table scans.
+   */
+  private void recalculateSnapshotDbPendingDeleteQuota(
+      Map<String, OmBucketInfo> nameBucketInfoMap,
+      Map<Long, OmBucketInfo> bucketById,
+      OMMetadataManager activeMetaManager,
+      Map<String, CountPair> snapshotDeletedKeyMap,
+      Map<String, CountPair> snapshotDeletedDirMap) throws IOException {
+    LOG.info("Starting recalculate snapshot pending-delete from snapshot DBs");
+    OMMetadataManager liveMetaManager = om.getMetadataManager();
+    SnapshotChainManager chain =
+        ((OmMetadataManagerImpl) liveMetaManager).getSnapshotChainManager();
+    OmSnapshotManager snapshotManager = om.getOmSnapshotManager();
+    Set<UUID> scannedSnapshotIds = new HashSet<>();
+
+    for (OmBucketInfo bucket : nameBucketInfoMap.values()) {
+      String snapshotPath = buildSnapshotPath(bucket.getVolumeName(), 
bucket.getBucketName());
+      LinkedHashMap<UUID, SnapshotChainInfo> pathChain;
+      try {
+        pathChain = chain.getSnapshotChainPath(snapshotPath);
+      } catch (IOException ex) {
+        throw new IOException("Failed to read snapshot chain for path " + 
snapshotPath, ex);
+      }
+      if (pathChain == null || pathChain.isEmpty()) {
+        continue;
+      }
+      for (UUID snapshotId : pathChain.keySet()) {
+        if (!scannedSnapshotIds.add(snapshotId)) {
+          continue;
+        }
+        SnapshotInfo snapshotInfo = loadActiveSnapshot(activeMetaManager, 
chain, snapshotId);
+        if (snapshotInfo == null) {
+          continue;
+        }
+        if (!snapshotInfo.getVolumeName().equals(bucket.getVolumeName())
+            || !snapshotInfo.getBucketName().equals(bucket.getBucketName())) {
+          continue;
+        }
+        if (!OmSnapshotManager.isSnapshotFlushedToDB(liveMetaManager, 
snapshotInfo)) {
+          LOG.warn("Skipping snapshot {} for quota repair: create txn not 
flushed to active DB",
+              snapshotInfo.getTableKey());
+          continue;
+        }
+        String sourceLabel = "snapshot DB " + snapshotInfo.getTableKey();
+        try (UncheckedAutoCloseableSupplier<OmSnapshot> snapshotRef =
+                 snapshotManager.getSnapshot(snapshotId)) {
+          scanDeletedTables(snapshotRef.get().getMetadataManager(), bucketById,
+              snapshotDeletedKeyMap, snapshotDeletedDirMap, sourceLabel);
+        }
+      }
+    }
+    LOG.info("Recalculate snapshot pending-delete from snapshot DBs completed, 
snapshots scanned: {}",
+        scannedSnapshotIds.size());
+  }
+
+  private static String buildSnapshotPath(String volumeName, String 
bucketName) {
+    return volumeName + OM_KEY_PREFIX + bucketName;
+  }
+
+  private static SnapshotInfo loadActiveSnapshot(
+      OMMetadataManager metadataManager,
+      SnapshotChainManager chain,
+      UUID snapshotId) throws IOException {
+    String tableKey = chain.getTableKey(snapshotId);
+    if (tableKey == null) {
+      LOG.warn("Snapshot id {} is not present in snapshot chain table-key 
map", snapshotId);
+      return null;
+    }
+    SnapshotInfo snapshotInfo = 
metadataManager.getSnapshotInfoTable().get(tableKey);
+    if (snapshotInfo == null) {
+      LOG.warn("Snapshot {} not found in snapshotInfoTable during quota 
repair", tableKey);
+      return null;
+    }
+    if (snapshotInfo.getSnapshotStatus() != SNAPSHOT_ACTIVE) {
+      return null;
+    }
+    return snapshotInfo;
+  }
+
+  private void scanDeletedTables(
+      OMMetadataManager metadataManager,
+      Map<Long, OmBucketInfo> bucketById,
+      Map<String, CountPair> snapshotDeletedKeyMap,
+      Map<String, CountPair> snapshotDeletedDirMap,
+      String sourceLabel) throws UncheckedIOException {
+    recalculateDeletedKeyUsages(metadataManager.getDeletedTable(), bucketById,
+        snapshotDeletedKeyMap, sourceLabel);
+    recalculateDeletedDirNamespace(metadataManager.getDeletedDirTable(),
+        snapshotDeletedDirMap, sourceLabel);
+  }
+
+  private void recalculateDeletedKeyUsages(
+      Table<String, RepeatedOmKeyInfo> deletedTable,
+      Map<Long, OmBucketInfo> bucketById,
+      Map<String, CountPair> snapshotCountByBucketNameKey,
+      String sourceLabel)
+      throws UncheckedIOException {
+    LOG.info("Starting recalculate snapshot usages from deletedTable ({})", 
sourceLabel);
+
+    int count = 0;
+    long startTime = Time.monotonicNow();
+    try (Table.KeyValueIterator<String, RepeatedOmKeyInfo> keyIter
+        = deletedTable.iterator()) {
+      while (keyIter.hasNext()) {
+        Table.KeyValue<String, RepeatedOmKeyInfo> kv = keyIter.next();
+        count++;
+        RepeatedOmKeyInfo val = kv.getValue();
+        OmBucketInfo bucket = bucketById.get(val.getBucketId());
+        if (bucket == null) {
+          continue;
+        }
+        String nameKey = buildNamePath(bucket.getVolumeName(), 
bucket.getBucketName());
+        CountPair usage = snapshotCountByBucketNameKey.get(nameKey);
+        if (usage == null) {
+          continue;
+        }
+        usage.incrSpace(val.getTotalSize().getRight());
+        usage.incrNamespace(val.getOmKeyInfoList().size());
+      }
+      LOG.info("Recalculate snapshot usages from deletedTable ({}) completed, 
count {} time {}ms",
+          sourceLabel, count, (Time.monotonicNow() - startTime));
+    } catch (IOException ex) {
+      throw new UncheckedIOException(ex);
+    }
+  }
+
+  private void recalculateDeletedDirNamespace(
+      Table<String, OmKeyInfo> deletedDirTable,
+      Map<String, CountPair> snapshotDirNsByIdPrefix,
+      String sourceLabel)
+      throws UncheckedIOException {
+    LOG.info("Starting recalculate snapshot namespace from 
deletedDirectoryTable ({})",
+        sourceLabel);
+
+    int count = 0;
+    long startTime = Time.monotonicNow();
+    try (Table.KeyValueIterator<String, OmKeyInfo> keyIter
+        = deletedDirTable.iterator()) {
+      while (keyIter.hasNext()) {
+        Table.KeyValue<String, OmKeyInfo> kv = keyIter.next();
+        count++;
+        String prefix = getVolumeBucketPrefix(kv.getKey());
+        CountPair usage = snapshotDirNsByIdPrefix.get(prefix);
+        if (usage != null) {
+          usage.incrNamespace(1L);
+        }
+      }
+      LOG.info(
+          "Recalculate snapshot namespace from deletedDirectoryTable ({}) 
completed, count {} time {}ms",
+          sourceLabel, count, (Time.monotonicNow() - startTime));
+    } catch (IOException ex) {
+      throw new UncheckedIOException(ex);
+    }
+  }
+
+  private static synchronized void mergeSnapshotDeletedTableCounts(
+      Map<String, OmBucketInfo> nameBucketInfoMap,
+      Map<String, CountPair> counts) {
+    for (Map.Entry<String, CountPair> entry : counts.entrySet()) {
+      OmBucketInfo bucket = nameBucketInfoMap.get(entry.getKey());
+      if (bucket != null) {
+        bucket.incrSnapshotUsedBytes(entry.getValue().getSpace());
+        bucket.incrSnapshotUsedNamespace(entry.getValue().getNamespace());
+      }
+    }
+  }
+
+  private static synchronized void mergeDeletedDirSnapshotNamespace(
+      Map<String, OmBucketInfo> idBucketInfoMap,
+      Map<String, CountPair> counts) {
+    for (Map.Entry<String, CountPair> entry : counts.entrySet()) {
+      OmBucketInfo bucket = idBucketInfoMap.get(entry.getKey());
+      if (bucket != null) {
+        bucket.incrSnapshotUsedNamespace(entry.getValue().getNamespace());
+      }
+    }
+  }
+
   private <VALUE> void recalculateUsages(
       Table<String, VALUE> table, Map<String, CountPair> prefixUsageMap,
       String strType, boolean haveValue) throws UncheckedIOException,
@@ -500,6 +743,12 @@ public void 
updateStatus(OzoneManagerProtocolProtos.QuotaRepairRequest.Builder b
         ConcurrentHashMap<String, Long> diffCountMap = new 
ConcurrentHashMap<>();
         diffCountMap.put("DiffUsedBytes", quotaCount.getDiffUsedBytes());
         diffCountMap.put("DiffUsedNamespace", 
quotaCount.getDiffUsedNamespace());
+        if (quotaCount.hasDiffSnapshotUsedBytes()) {
+          diffCountMap.put("DiffSnapshotUsedBytes", 
quotaCount.getDiffSnapshotUsedBytes());
+        }
+        if (quotaCount.hasDiffSnapshotUsedNamespace()) {
+          diffCountMap.put("DiffSnapshotUsedNamespace", 
quotaCount.getDiffSnapshotUsedNamespace());
+        }
         bucketCountDiffMap.put(bucketKey, diffCountMap);
       }
     }
diff --git 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestQuotaRepairTask.java
 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestQuotaRepairTask.java
index de950e8a5a1..cf8ddfbba56 100644
--- 
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestQuotaRepairTask.java
+++ 
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestQuotaRepairTask.java
@@ -20,6 +20,7 @@
 import static 
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.ONE;
 import static 
org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyLong;
@@ -29,6 +30,7 @@
 
 import java.io.IOException;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
 import org.apache.hadoop.hdds.client.RatisReplicationConfig;
 import org.apache.hadoop.hdds.utils.db.BatchOperation;
@@ -38,6 +40,7 @@
 import org.apache.hadoop.ozone.om.helpers.OmBucketInfo;
 import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
 import org.apache.hadoop.ozone.om.helpers.OmVolumeArgs;
+import org.apache.hadoop.ozone.om.helpers.RepeatedOmKeyInfo;
 import org.apache.hadoop.ozone.om.ratis.OzoneManagerRatisServer;
 import org.apache.hadoop.ozone.om.request.OMRequestTestUtils;
 import org.apache.hadoop.ozone.om.request.key.TestOMKeyRequest;
@@ -47,12 +50,21 @@
 import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
 import org.apache.hadoop.util.Time;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
 
 /**
  * Test class for quota repair.
  */
+@Timeout(120)
 public class TestQuotaRepairTask extends TestOMKeyRequest {
 
+  /** Seconds; must match {@link Timeout} on this class. */
+  private static final int REPAIR_TEST_TIMEOUT_SECONDS = 120;
+
+  private static Boolean awaitRepair(CompletableFuture<Boolean> repair) throws 
Exception {
+    return repair.get(REPAIR_TEST_TIMEOUT_SECONDS, TimeUnit.SECONDS);
+  }
+
   @Test
   public void testQuotaRepair() throws Exception {
     OzoneManagerProtocolProtos.OMResponse respMock = 
mock(OzoneManagerProtocolProtos.OMResponse.class);
@@ -110,7 +122,7 @@ public void testQuotaRepair() throws Exception {
     
     QuotaRepairTask quotaRepairTask = new QuotaRepairTask(ozoneManager);
     CompletableFuture<Boolean> repair = quotaRepairTask.repair();
-    Boolean repairStatus = repair.get();
+    Boolean repairStatus = awaitRepair(repair);
     assertTrue(repairStatus);
 
     OMQuotaRepairRequest omQuotaRepairRequest = new 
OMQuotaRepairRequest(ref.get());
@@ -170,7 +182,7 @@ public void testQuotaRepairForOldVersionVolumeBucket() 
throws Exception {
 
     QuotaRepairTask quotaRepairTask = new QuotaRepairTask(ozoneManager);
     CompletableFuture<Boolean> repair = quotaRepairTask.repair();
-    Boolean repairStatus = repair.get();
+    Boolean repairStatus = awaitRepair(repair);
     assertTrue(repairStatus);
 
     OMQuotaRepairRequest omQuotaRepairRequest = new 
OMQuotaRepairRequest(ref.get());
@@ -187,6 +199,132 @@ public void testQuotaRepairForOldVersionVolumeBucket() 
throws Exception {
     assertEquals(-1, volArgsVerify.getQuotaInNamespace());
   }
 
+  @Test
+  public void testQuotaRepairDeletedTableSnapshotQuota() throws Exception {
+    OzoneManagerProtocolProtos.OMResponse respMock = 
mock(OzoneManagerProtocolProtos.OMResponse.class);
+    when(respMock.getSuccess()).thenReturn(true);
+    OzoneManagerRatisServer ratisServerMock = 
mock(OzoneManagerRatisServer.class);
+    AtomicReference<OzoneManagerProtocolProtos.OMRequest> ref = new 
AtomicReference<>();
+    doAnswer(invocation -> {
+      ref.set(invocation.getArgument(0, 
OzoneManagerProtocolProtos.OMRequest.class));
+      return respMock;
+    }).when(ratisServerMock).submitRequest(any(), any(), anyLong());
+    when(ozoneManager.getOmRatisServer()).thenReturn(ratisServerMock);
+
+    OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+        omMetadataManager, BucketLayout.OBJECT_STORE);
+
+    String keyName = "/user/snapKey";
+    OMRequestTestUtils.addKeyToTableAndCache(volumeName, bucketName,
+        keyName, -1, RatisReplicationConfig.getInstance(THREE), 1L, 
omMetadataManager);
+
+    String ozoneKey = omMetadataManager.getOzoneKey(volumeName, bucketName, 
keyName);
+    OmBucketInfo bucketInfo = omMetadataManager.getBucketTable().get(
+        omMetadataManager.getBucketKey(volumeName, bucketName));
+    long bucketObjId = bucketInfo.getObjectID();
+
+    OMRequestTestUtils.deleteKey(ozoneKey, bucketObjId, omMetadataManager, 2L);
+
+    RepeatedOmKeyInfo deletedEntry = 
omMetadataManager.getDeletedTable().get(ozoneKey);
+    long expectedSnapNs = deletedEntry.getOmKeyInfoList().size();
+
+    bucketInfo = omMetadataManager.getBucketTable().get(
+        omMetadataManager.getBucketKey(volumeName, bucketName));
+    OmBucketInfo corruptedSnapshot = bucketInfo.toBuilder()
+        .setSnapshotUsedBytes(7L)
+        .setSnapshotUsedNamespace(99L)
+        .build();
+    String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName);
+    omMetadataManager.getBucketTable().put(bucketKey, corruptedSnapshot);
+    omMetadataManager.getBucketTable().addCacheEntry(
+        new CacheKey<>(bucketKey), CacheValue.get(3L, corruptedSnapshot));
+
+    QuotaRepairTask quotaRepairTask = new QuotaRepairTask(ozoneManager);
+    CompletableFuture<Boolean> repair = quotaRepairTask.repair();
+    assertTrue(awaitRepair(repair));
+
+    OMQuotaRepairRequest omQuotaRepairRequest = new 
OMQuotaRepairRequest(ref.get());
+    OMClientResponse omClientResponse = 
omQuotaRepairRequest.validateAndUpdateCache(ozoneManager, 1);
+    BatchOperation batchOperation = 
omMetadataManager.getStore().initBatchOperation();
+    ((OMQuotaRepairResponse) omClientResponse).addToDBBatch(omMetadataManager, 
batchOperation);
+    omMetadataManager.getStore().commitBatchOperation(batchOperation);
+
+    OmBucketInfo repaired = omMetadataManager.getBucketTable().get(bucketKey);
+    assertEquals(0, repaired.getUsedBytes());
+    assertEquals(0, repaired.getUsedNamespace());
+    assertEquals(expectedSnapNs, repaired.getSnapshotUsedNamespace());
+    assertTrue(repaired.getSnapshotUsedBytes() > 0,
+        "Snapshot pending-delete bytes must be recomputed from deletedTable");
+  }
+
+  @Test
+  public void testQuotaRepairSnapshotDbDeletedTableQuota() throws Exception {
+    OzoneManagerProtocolProtos.OMResponse respMock = 
mock(OzoneManagerProtocolProtos.OMResponse.class);
+    when(respMock.getSuccess()).thenReturn(true);
+    OzoneManagerRatisServer ratisServerMock = 
mock(OzoneManagerRatisServer.class);
+    AtomicReference<OzoneManagerProtocolProtos.OMRequest> ref = new 
AtomicReference<>();
+    doAnswer(invocation -> {
+      ref.set(invocation.getArgument(0, 
OzoneManagerProtocolProtos.OMRequest.class));
+      return respMock;
+    }).when(ratisServerMock).submitRequest(any(), any(), anyLong());
+    when(ozoneManager.getOmRatisServer()).thenReturn(ratisServerMock);
+
+    OMRequestTestUtils.addVolumeAndBucketToDB(volumeName, bucketName,
+        omMetadataManager, BucketLayout.OBJECT_STORE);
+
+    String keyName = "/user/snapKey";
+    OMRequestTestUtils.addKeyToTableAndCache(volumeName, bucketName,
+        keyName, -1, RatisReplicationConfig.getInstance(THREE), 1L, 
omMetadataManager);
+
+    String ozoneKey = omMetadataManager.getOzoneKey(volumeName, bucketName, 
keyName);
+    OmKeyInfo omKeyInfo = 
omMetadataManager.getKeyTable(BucketLayout.OBJECT_STORE).get(ozoneKey);
+    long keyBytes = omKeyInfo.getReplicatedSize();
+
+    OmBucketInfo bucketInfo = omMetadataManager.getBucketTable().get(
+        omMetadataManager.getBucketKey(volumeName, bucketName));
+    OMRequestTestUtils.deleteKey(ozoneKey, bucketInfo.getObjectID(), 
omMetadataManager, 2L);
+
+    String bucketKey = omMetadataManager.getBucketKey(volumeName, bucketName);
+    OmBucketInfo afterDelete = bucketInfo.toBuilder()
+        .setUsedBytes(0)
+        .setUsedNamespace(0)
+        .setSnapshotUsedBytes(keyBytes)
+        .setSnapshotUsedNamespace(1)
+        .build();
+    omMetadataManager.getBucketTable().put(bucketKey, afterDelete);
+
+    when(ozoneManager.getDefaultReplicationConfig())
+        .thenReturn(RatisReplicationConfig.getInstance(THREE));
+    createSnapshot("snap1");
+
+    assertNull(omMetadataManager.getDeletedTable().get(ozoneKey),
+        "Deleted key should move out of active deletedTable after snapshot");
+    assertEquals(0, 
omMetadataManager.countRowsInTable(omMetadataManager.getDeletedTable()));
+
+    OmBucketInfo corrupted = afterDelete.toBuilder()
+        .setSnapshotUsedBytes(7L)
+        .build();
+    omMetadataManager.getBucketTable().put(bucketKey, corrupted);
+    omMetadataManager.getBucketTable().addCacheEntry(
+        new CacheKey<>(bucketKey), CacheValue.get(3L, corrupted));
+
+    QuotaRepairTask quotaRepairTask = new QuotaRepairTask(ozoneManager);
+    CompletableFuture<Boolean> repair = quotaRepairTask.repair();
+    assertTrue(awaitRepair(repair));
+
+    OMQuotaRepairRequest omQuotaRepairRequest = new 
OMQuotaRepairRequest(ref.get());
+    OMClientResponse omClientResponse = 
omQuotaRepairRequest.validateAndUpdateCache(ozoneManager, 1);
+    BatchOperation batchOperation = 
omMetadataManager.getStore().initBatchOperation();
+    ((OMQuotaRepairResponse) omClientResponse).addToDBBatch(omMetadataManager, 
batchOperation);
+    omMetadataManager.getStore().commitBatchOperation(batchOperation);
+
+    OmBucketInfo repaired = omMetadataManager.getBucketTable().get(bucketKey);
+    assertEquals(0, repaired.getUsedBytes());
+    assertEquals(0, repaired.getUsedNamespace());
+    assertEquals(keyBytes, repaired.getSnapshotUsedBytes());
+    assertEquals(1, repaired.getSnapshotUsedNamespace());
+  }
+
   private void zeroOutBucketUsedBytes(String volumeName, String bucketName,
                                       long trxnLogIndex)
       throws IOException {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to