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


##########
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:
   Added GenericTestUtils.waitFor before checking the completed-task count, 
matching the other tests here. The test expects one completed task after 
partial rejection and zero for a stopped executor.



-- 
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