rich7420 commented on code in PR #11009:
URL: https://github.com/apache/ozone/pull/11009#discussion_r4129051178


##########
hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java:
##########
@@ -670,44 +685,107 @@ void processDeletedDirsForStore(SnapshotInfo 
currentSnapshotInfo, KeyManager key
         UUID expectedPreviousSnapshotId = currentSnapshotInfo == null ?
             snapshotChainManager.getLatestGlobalSnapshotId() :
             SnapshotUtils.getPreviousSnapshotId(currentSnapshotInfo, 
snapshotChainManager);
-        Map<UUID, Pair<Long, Long>> exclusiveSizeMap = Maps.newConcurrentMap();
+        processedAllDeletedDirs = runDeletionWorkers(currentSnapshotInfo, 
keyManager, dirSupplier, currentSnapshot,
+            expectedPreviousSnapshotId, exclusiveSizeMap, rnCnt, remainNum);
+      }
+
+      // If AOS or all directories have been processed for snapshot, update 
snapshot size delta and deep clean flag
+      // if it is a snapshot. All snapshot DB iterators and handles have been 
closed before this synchronous request.
+      if (processedAllDeletedDirs) {
+        List<OzoneManagerProtocolProtos.SetSnapshotPropertyRequest> 
setSnapshotPropertyRequests = new ArrayList<>();
+
+        for (Map.Entry<UUID, Pair<Long, Long>> entry : 
exclusiveSizeMap.entrySet()) {
+          UUID snapshotID = entry.getKey();
+          long exclusiveSize = entry.getValue().getLeft();
+          long exclusiveReplicatedSize = entry.getValue().getRight();
+          
setSnapshotPropertyRequests.add(getSetSnapshotRequestUpdatingExclusiveSize(
+              exclusiveSize, exclusiveReplicatedSize, snapshotID));
+        }
+
+        // Updating directory deep clean flag of snapshot.
+        if (currentSnapshotInfo != null) {
+          
setSnapshotPropertyRequests.add(OzoneManagerProtocolProtos.SetSnapshotPropertyRequest.newBuilder()
+              .setSnapshotKey(snapshotTableKey)
+              .setDeepCleanedDeletedDir(true)
+              .build());
+        }
+        submitSetSnapshotRequests(setSnapshotPropertyRequests);
+      }
+    }
 
-        CompletableFuture<Boolean> processedAllDeletedDirs = 
CompletableFuture.completedFuture(true);
-        final int parallelThreads = numberOfParallelThreadsPerStore.get();
-        for (int i = 0; i < parallelThreads; i++) {
+    /**
+     * Waits for snapshot workers to close DB handles before opening the 
submission gate.
+     * The workers retain their GC locks through submission; the coordinator 
closes its own iterator and handle.
+     * AOS has no task-owned snapshot handle, so its workers can close their 
handles and submit independently.
+     */
+    @SuppressWarnings("checkstyle:ParameterNumber")
+    private boolean runDeletionWorkers(SnapshotInfo currentSnapshotInfo, 
KeyManager keyManager,
+        DeletedDirSupplier dirSupplier, 
UncheckedAutoCloseableSupplier<OmSnapshot> currentSnapshot,
+        UUID expectedPreviousSnapshotId, Map<UUID, Pair<Long, Long>> 
exclusiveSizeMap, long rnCnt, int remainNum)
+        throws ExecutionException, InterruptedException {
+      int parallelThreads = numberOfParallelThreadsPerStore.get();
+      CountDownLatch snapshotDbHandlesClosed = currentSnapshotInfo == null ? 
null : new CountDownLatch(parallelThreads);
+      CountDownLatch submitRequests = new CountDownLatch(1);
+      CompletableFuture<Boolean> workersCompleted = 
CompletableFuture.completedFuture(true);
+      for (int i = 0; i < parallelThreads; i++) {
+        AtomicBoolean workerReady = new AtomicBoolean();
+        try {
           CompletableFuture<Boolean> future = CompletableFuture.supplyAsync(() 
-> {
             try {
               return processDeletedDirectories(currentSnapshotInfo, 
keyManager, dirSupplier,
-                  expectedPreviousSnapshotId, exclusiveSizeMap, rnCnt, 
remainNum);
+                  expectedPreviousSnapshotId, exclusiveSizeMap, rnCnt, 
remainNum, () -> {
+                    if (snapshotDbHandlesClosed != null) {
+                      signalWorkerReady(workerReady, snapshotDbHandlesClosed);
+                      if (awaitUninterruptibly(submitRequests)) {
+                        Thread.currentThread().interrupt();
+                      }
+                    }
+                  });
             } catch (Throwable e) {
               return false;
+            } finally {
+              if (snapshotDbHandlesClosed != null) {
+                signalWorkerReady(workerReady, snapshotDbHandlesClosed);
+              }
             }
           }, isThreadPoolActive(deletionThreadPool) ? deletionThreadPool : 
ForkJoinPool.commonPool());

Review Comment:
   This fallback can deadlock the snapshot coordinator during shutdown. With 
three workers and common-pool parallelism two, two workers wait at the 
submission gate while the third remains queued. I reproduced this with a 
stopped DDS executor and a real snapshot. Please submit directly to the DDS 
executor and use the rejection handling below; the same test passes with that 
change.



##########
hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestDirectoryDeletingService.java:
##########
@@ -383,4 +413,254 @@ void testPurgeDirectoriesGroupedByBucketPerTransaction() 
throws Exception {
         bucketIdA, bucketIdB);
   }
 
+  /**
+   * Regression test for HDDS-16164. Snapshot workers must wait for the 
task-owned handle to close before submitting.
+   * AOS workers own their handles and can finish scanning while a sibling 
waits for double-buffer capacity.
+   */
+  @ParameterizedTest(name = "activeStore={0}")
+  @ValueSource(booleans = {false, true})
+  @DisplayName("DirectoryDeletingService must not deadlock at the unflushed 
transaction limit")
+  void testNoUnflushedTransactionDeadlock(boolean activeStore) throws 
Exception {
+    int threadCount = 3;
+    OzoneConfiguration conf = createConfAndInitValues(threadCount);
+    conf.setInt(OMConfigKeys.OZONE_OM_UNFLUSHED_TRANSACTION_MAX_COUNT, 1);
+    conf.setInt(OZONE_PATH_DELETING_LIMIT_PER_TASK, 1);
+    conf.setInt(OZONE_MANAGER_STRIPED_LOCK_SIZE_PREFIX + "snapshot_db_lock", 
1);
+    conf.setTimeDuration(OZONE_BLOCK_DELETING_SERVICE_INTERVAL, 1, 
TimeUnit.HOURS);
+    conf.setTimeDuration(OZONE_SNAPSHOT_DELETING_SERVICE_INTERVAL, 1, 
TimeUnit.HOURS);
+    conf.setTimeDuration(OZONE_DIR_DELETING_SERVICE_INTERVAL, 1, 
TimeUnit.HOURS);
+
+    OmTestManagers managers = new OmTestManagers(conf);
+    om = managers.getOzoneManager();
+    KeyManager keyManager = managers.getKeyManager();
+    OMMetadataManager metadataManager = managers.getMetadataManager();
+    OzoneManagerProtocol writeClient = managers.getWriteClient();
+    keyManager.getDirDeletingService().shutdown();
+
+    OMRequestTestUtils.addVolumeToOM(metadataManager,
+        
OmVolumeArgs.newBuilder().setOwnerName("o").setAdminName("a").setVolume(volumeName).build());
+    OMRequestTestUtils.addBucketToOM(metadataManager,
+        
OmBucketInfo.newBuilder().setVolumeName(volumeName).setBucketName(bucketName)
+            
.setObjectID(1).setBucketLayout(BucketLayout.FILE_SYSTEM_OPTIMIZED).build());
+
+    String purgeSnapshotName = "purge";
+    writeClient.createSnapshot(volumeName, bucketName, purgeSnapshotName);
+    String previousSnapshotName = "previous";
+    writeClient.createSnapshot(volumeName, bucketName, previousSnapshotName);
+    String deepCleanSnapshotName = "deep-clean";
+    if (activeStore) {
+      writeClient.createSnapshot(volumeName, bucketName, 
deepCleanSnapshotName);
+    }
+    om.awaitDoubleBufferFlush();
+
+    addDeletedDirectories(metadataManager, threadCount);
+    if (!activeStore) {
+      writeClient.createSnapshot(volumeName, bucketName, 
deepCleanSnapshotName);
+    }
+    om.awaitDoubleBufferFlush();
+
+    String deepCleanKey = SnapshotInfo.getTableKey(volumeName, bucketName, 
deepCleanSnapshotName);
+    SnapshotInfo deepCleanInfo = 
metadataManager.getSnapshotInfoTable().get(deepCleanKey);
+    assertNotNull(deepCleanInfo);
+    awaitSnapshotFlushed(metadataManager, deepCleanInfo);
+
+    writeClient.deleteSnapshot(volumeName, bucketName, purgeSnapshotName);
+    om.awaitDoubleBufferFlush();
+    String purgeKey = SnapshotInfo.getTableKey(volumeName, bucketName, 
purgeSnapshotName);
+    assertNotNull(metadataManager.getSnapshotInfoTable().get(purgeKey));
+
+    CountDownLatch ddsAtSubmit = new CountDownLatch(1);
+    CountDownLatch ddsMayProceed = new CountDownLatch(1);
+    DirectoryDeletingService service = Mockito.spy(new 
DirectoryDeletingService(1, TimeUnit.HOURS, 1,
+        om, conf, threadCount, true) {
+      @Override
+      protected OzoneManagerProtocolProtos.OMResponse submitRequest(
+          OzoneManagerProtocolProtos.OMRequest omRequest) throws 
ServiceException {
+        assertThat(omRequest.getCmdType()).isIn(
+            OzoneManagerProtocolProtos.Type.PurgeDirectories,
+            OzoneManagerProtocolProtos.Type.SetSnapshotProperty);
+        ddsAtSubmit.countDown();
+        try {
+          ddsMayProceed.await();
+        } catch (InterruptedException e) {
+          Thread.currentThread().interrupt();
+          throw new ServiceException(e);
+        }
+        return super.submitRequest(omRequest);
+      }
+    });
+
+    CountDownLatch workersPrepared = new CountDownLatch(threadCount);
+    CountDownLatch finishScanning = new CountDownLatch(1);
+    AtomicBoolean pauseScan = new AtomicBoolean();
+    Mockito.doAnswer(invocation -> {
+      workersPrepared.countDown();
+      if (pauseScan.compareAndSet(false, true)) {
+        // Keep this worker's previous-snapshot handle open while the other 
workers finish scanning.
+        assertThat(metadataManager.getLock().getReadHoldCount(SNAPSHOT_DB_LOCK,
+            deepCleanInfo.getSnapshotId().toString())).isPositive();
+        assertTrue(finishScanning.await(30, TimeUnit.SECONDS));
+      }
+      return invocation.callRealMethod();
+    }).when(service).optimizeDirDeletesAndSubmitRequest(Mockito.anyLong(), 
Mockito.anyLong(), Mockito.anyLong(),
+        Mockito.anyList(), Mockito.anyList(), Mockito.any(), 
Mockito.anyLong(), Mockito.any(), Mockito.any(),
+        Mockito.any(), Mockito.anyMap(), Mockito.any(), Mockito.anyLong(), 
Mockito.any(), Mockito.any());
+
+    OzoneManagerDoubleBuffer doubleBuffer =
+        
om.getOmRatisServer().getOmStateMachine().getOzoneManagerDoubleBuffer();
+    Semaphore unflushedTransactions =
+        (Semaphore) HddsWhiteboxTestUtils.getInternalState(doubleBuffer, 
"unFlushedTransactions");
+    ThreadPoolExecutor deletionThreadPool = service.getDeletionThreadPool();
+    long completedTaskCountBefore = deletionThreadPool.getCompletedTaskCount();
+    ExecutorService executor = Executors.newFixedThreadPool(2);
+    try {
+      Future<?> ddsFuture = service.getExecutorService().submit(
+          () -> service.new DirDeletingTask(activeStore ? null : 
deepCleanInfo.getSnapshotId()).call());
+      assertTrue(workersPrepared.await(30, TimeUnit.SECONDS), "Expected all 
configured workers to scan in parallel");
+      if (!activeStore) {
+        assertThat(ddsAtSubmit.await(200, TimeUnit.MILLISECONDS))
+            .as("snapshot workers must wait for the task-owned snapshot handle 
to close").isFalse();
+        finishScanning.countDown();
+      }
+      assertTrue(ddsAtSubmit.await(30, TimeUnit.SECONDS), "DirDeletingTask 
never reached submitRequest");
+
+      CountDownLatch nextStoreStarted = new CountDownLatch(1);
+      Future<?> nextStore = 
service.getExecutorService().submit(nextStoreStarted::countDown);
+      assertThat(nextStoreStarted.await(200, TimeUnit.MILLISECONDS))
+          .as("the next store must not open DB handles while workers are 
submitting requests").isFalse();
+
+      AtomicBoolean purgeDone = new AtomicBoolean();
+      OzoneManagerProtocolProtos.OMRequest purgeRequest = 
OzoneManagerProtocolProtos.OMRequest.newBuilder()
+          .setCmdType(OzoneManagerProtocolProtos.Type.SnapshotPurge)
+          
.setSnapshotPurgeRequest(OzoneManagerProtocolProtos.SnapshotPurgeRequest.newBuilder()
+              .addSnapshotDBKeys(purgeKey))
+          .setClientId(ClientId.randomId().toString())
+          .build();
+      Future<?> purgeFuture = executor.submit(() -> {
+        OzoneManagerRatisUtils.submitRequest(om, purgeRequest, 
ClientId.randomId(), 1L);
+        purgeDone.set(true);
+        return null;
+      });
+      GenericTestUtils.waitFor(
+          () -> unflushedTransactions.availablePermits() == 0 || 
purgeDone.get(), 50, 30000);
+
+      ddsMayProceed.countDown();
+      if (activeStore) {
+        GenericTestUtils.waitFor(unflushedTransactions::hasQueuedThreads, 50, 
10000);
+      }
+      // An AOS worker can finish scanning and close its handle even while a 
sibling waits for flush capacity.
+      finishScanning.countDown();
+      try {
+        ddsFuture.get(60, TimeUnit.SECONDS);
+      } catch (TimeoutException e) {
+        fail("DirectoryDeletingService retained a snapshot DB read lock across 
its synchronous OM request");
+      }
+      purgeFuture.get(60, TimeUnit.SECONDS);
+      nextStore.get(10, TimeUnit.SECONDS);
+      om.awaitDoubleBufferFlush();
+      GenericTestUtils.waitFor(() -> 
deletionThreadPool.getCompletedTaskCount() ==
+          completedTaskCountBefore + threadCount, 100, 5000);
+      assertThat(deletionThreadPool.getCompletedTaskCount() - 
completedTaskCountBefore)
+          .as("store processing should use all configured deletion 
workers").isEqualTo(threadCount);
+      assertThat(service.getDeletedDirsCount()).isEqualTo(threadCount);
+    } finally {
+      finishScanning.countDown();
+      ddsMayProceed.countDown();
+      executor.shutdownNow();
+      service.shutdown();
+    }
+  }
+
+  @Test
+  @DisplayName("DirectoryDeletingService releases snapshot workers when 
submission is rejected")
+  void testRejectedSnapshotWorkerSubmission() throws Exception {
+    OzoneConfiguration conf = createConfAndInitValues(2);
+    OmTestManagers managers = new OmTestManagers(conf);
+    om = managers.getOzoneManager();
+    KeyManager keyManager = managers.getKeyManager();
+    OMMetadataManager metadataManager = managers.getMetadataManager();
+    OzoneManagerProtocol writeClient = managers.getWriteClient();
+    keyManager.getDirDeletingService().shutdown();
+
+    OMRequestTestUtils.addVolumeToOM(metadataManager,
+        
OmVolumeArgs.newBuilder().setOwnerName("o").setAdminName("a").setVolume(volumeName).build());
+    OMRequestTestUtils.addBucketToOM(metadataManager,
+        
OmBucketInfo.newBuilder().setVolumeName(volumeName).setBucketName(bucketName)
+            
.setObjectID(1).setBucketLayout(BucketLayout.FILE_SYSTEM_OPTIMIZED).build());
+    String snapshotName = "snapshot";
+    writeClient.createSnapshot(volumeName, bucketName, snapshotName);
+    om.awaitDoubleBufferFlush();
+
+    String snapshotKey = SnapshotInfo.getTableKey(volumeName, bucketName, 
snapshotName);
+    SnapshotInfo snapshotInfo = 
metadataManager.getSnapshotInfoTable().get(snapshotKey);
+    assertNotNull(snapshotInfo);
+    awaitSnapshotFlushed(metadataManager, snapshotInfo);
+
+    DirectoryDeletingService service = new DirectoryDeletingService(1, 
TimeUnit.HOURS, 1,
+        om, conf, 2, true);
+    RejectAfterFirstExecutor rejectionExecutor = new 
RejectAfterFirstExecutor();
+    HddsWhiteboxTestUtils.setInternalState(service, "deletionThreadPool", 
rejectionExecutor);
+    ExecutorService taskExecutor = Executors.newSingleThreadExecutor();
+    try {
+      Future<?> processFuture = taskExecutor.submit(() -> {
+        try (UncheckedAutoCloseableSupplier<OmSnapshot> snapshot = 
om.getOmSnapshotManager().getActiveSnapshot(
+            volumeName, bucketName, snapshotName)) {
+          service.new 
DirDeletingTask(snapshotInfo.getSnapshotId()).processDeletedDirsForStore(snapshotInfo,
+              snapshot.get().getKeyManager(), snapshot, 1, 1);
+        }
+        return null;
+      });
+      processFuture.get(10, TimeUnit.SECONDS);
+      assertThat(rejectionExecutor.getCompletedTaskCount()).isEqualTo(1);

Review Comment:
   Please wait for the completed-task count with `GenericTestUtils.waitFor`, as 
the other tests here do. The future can complete before the executor increments 
the counter. Holding the worker in `afterExecute` reproduces `expected: 1, but 
was: 0`.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to