CRZbulabula commented on code in PR #18285:
URL: https://github.com/apache/iotdb/pull/18285#discussion_r3631211996


##########
iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java:
##########
@@ -305,34 +307,40 @@ public TSStatus deleteDatabases(
     for (final TDatabaseSchema databaseSchema : deleteSgSchemaList) {
       final String database = databaseSchema.getName();
       boolean hasOverlappedTask = false;
-      synchronized (this) {
-        while (executor.isRunning()
-            && System.currentTimeMillis() - startCheckTimeForProcedures < 
PROCEDURE_WAIT_TIME_OUT) {
-          final Pair<Long, Boolean> procedureIdDuplicatePair =
-              checkDuplicateTableTask(
-                  database, null, null, null, null, 
ProcedureType.DELETE_DATABASE_PROCEDURE);
-          hasOverlappedTask = procedureIdDuplicatePair.getRight();
-
-          if (Boolean.FALSE.equals(procedureIdDuplicatePair.getRight())) {
-            DeleteDatabaseProcedure procedure =
-                new DeleteDatabaseProcedure(databaseSchema, isGeneratedByPipe);
-            this.executor.submitProcedure(procedure);
-            procedures.add(procedure);
-            break;
+      while (executor.isRunning()
+          && System.currentTimeMillis() - startCheckTimeForProcedures < 
PROCEDURE_WAIT_TIME_OUT) {
+        try (final DatabaseLock ignored = 
acquireDatabaseLifecycleLock(database)) {
+          synchronized (this) {

Review Comment:
   Fixed in 14aba5c9f6. deleteDatabases now uses an interruptible timed lock 
acquisition with the remaining monotonic deadline, and returns the existing 
overlap/retry status if the lifecycle lock cannot be acquired before that 
deadline. Added testTimedRequestLockAcquisitionHonorsTimeout for the bounded 
wait.



##########
iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/scheduler/DatabaseLifecycleLockManager.java:
##########
@@ -0,0 +1,279 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.confignode.procedure.scheduler;
+
+import org.apache.iotdb.confignode.procedure.Procedure;
+
+import java.util.ArrayDeque;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.TreeSet;
+import java.util.concurrent.locks.Condition;
+import java.util.concurrent.locks.ReentrantLock;
+
+/**
+ * The single source of database lifecycle locks in a ConfigNode.
+ *
+ * <p>A lock is keyed by the exact database name. Both short-lived manager 
requests and procedures
+ * use the same ownership table, so database creation, Region creation, 
maintenance retries, and
+ * database deletion cannot bypass one another. Procedure ownership is 
identified by procedure id
+ * rather than worker thread because a procedure can resume on a different 
executor thread.
+ */
+public class DatabaseLifecycleLockManager {
+
+  private final ProcedureScheduler scheduler;
+  private final ReentrantLock stateLock = new ReentrantLock(true);
+  private final Condition lockReleased = stateLock.newCondition();
+  private final Map<String, DatabaseLockState> lockStateMap = new HashMap<>();
+
+  public DatabaseLifecycleLockManager(final ProcedureScheduler scheduler) {
+    this.scheduler = scheduler;
+  }
+
+  /** Acquire database locks for a synchronous manager request. */
+  public DatabaseLock acquireLocks(final Set<String> databaseNames) {
+    final List<String> orderedDatabases = orderedDatabases(databaseNames);
+    final Thread owner = Thread.currentThread();
+    stateLock.lock();
+    try {
+      while (!canAcquireRequestLocks(owner, orderedDatabases)) {
+        lockReleased.awaitUninterruptibly();
+      }
+      orderedDatabases.forEach(
+          database ->
+              lockStateMap
+                  .computeIfAbsent(database, ignored -> new 
DatabaseLockState())
+                  .acquireRequestLock(owner));
+      return new DatabaseLock(this, orderedDatabases, owner);
+    } finally {
+      stateLock.unlock();
+    }
+  }
+
+  /**
+   * Atomically tries to lock all databases for a procedure.
+   *
+   * @return the first database whose lock is unavailable, or null when all 
locks are acquired
+   */
+  public String tryLock(final Procedure<?> procedure, final Set<String> 
databaseNames) {
+    stateLock.lock();
+    try {
+      final List<String> acquiredDatabases = new ArrayList<>();
+      for (final String database : orderedDatabases(databaseNames)) {
+        final DatabaseLockState lockState =
+            lockStateMap.computeIfAbsent(database, ignored -> new 
DatabaseLockState());
+        if (!lockState.canAcquireProcedureLock(procedure)) {
+          acquiredDatabases.forEach(
+              acquiredDatabase -> releaseProcedureLock(procedure, 
acquiredDatabase));
+          return database;
+        }
+        if (lockState.acquireProcedureLock(procedure)) {
+          acquiredDatabases.add(database);
+        }
+      }
+      return null;
+    } finally {
+      stateLock.unlock();
+    }
+  }
+
+  public void waitProcedure(final Procedure<?> procedure, final String 
databaseName) {
+    stateLock.lock();
+    try {
+      final DatabaseLockState lockState =
+          lockStateMap.computeIfAbsent(databaseName, ignored -> new 
DatabaseLockState());
+      if (lockState.isUnlocked()) {
+        scheduler.addFront(procedure);
+        removeIfIdle(databaseName, lockState);
+      } else {
+        lockState.waitProcedure(procedure);
+      }
+    } finally {
+      stateLock.unlock();
+    }
+  }
+
+  public void releaseLocks(final Procedure<?> procedure, final Set<String> 
databaseNames) {
+    stateLock.lock();
+    try {
+      orderedDatabases(databaseNames)
+          .forEach(database -> releaseProcedureLock(procedure, database));
+    } finally {
+      stateLock.unlock();
+    }
+  }
+
+  private boolean canAcquireRequestLocks(final Thread owner, final 
List<String> orderedDatabases) {
+    for (final String database : orderedDatabases) {
+      final DatabaseLockState lockState = lockStateMap.get(database);
+      if (lockState != null && !lockState.canAcquireRequestLock(owner)) {
+        return false;
+      }
+    }
+    return true;
+  }
+
+  private void releaseRequestLocks(final List<String> orderedDatabases, final 
Thread requestOwner) {
+    stateLock.lock();
+    try {
+      for (final String database : orderedDatabases) {
+        final DatabaseLockState lockState = lockStateMap.get(database);
+        if (lockState != null && lockState.releaseRequestLock(requestOwner)) {
+          wakeWaiters(lockState);
+          removeIfIdle(database, lockState);
+        }
+      }
+    } finally {
+      stateLock.unlock();
+    }
+  }
+
+  private void releaseProcedureLock(final Procedure<?> procedure, final String 
database) {
+    final DatabaseLockState lockState = lockStateMap.get(database);
+    if (lockState != null && lockState.releaseProcedureLock(procedure)) {
+      wakeWaiters(lockState);
+      removeIfIdle(database, lockState);
+    }
+  }
+
+  private void wakeWaiters(final DatabaseLockState lockState) {
+    lockState.wakeWaitingProcedures(scheduler);
+    lockReleased.signalAll();
+  }

Review Comment:
   Fixed in 14aba5c9f6. The lock state now retains the waiting queue, wakes 
only its head, and does not signal synchronous request waiters while a 
Procedure is pending. The head priority is consumed atomically when that 
Procedure re-enters tryLock, so a request cannot overtake it. Added 
testWaitingProcedureCannotBeOvertakenByRequest to cover the release/reacquire 
window.



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

Reply via email to