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

yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/main by this push:
     new 242a736c8d [#13435] fix(core): re-disable a user-disabled metalake 
when force drop fails (#13436)
242a736c8d is described below

commit 242a736c8d69d9851aea8b4517ce58b86e59521a
Author: YangJie <[email protected]>
AuthorDate: Mon Sep 28 05:45:46 2026 -0400

    [#13435] fix(core): re-disable a user-disabled metalake when force drop 
fails (#13436)
    
    ### What changes were proposed in this pull request?
    
    `dropCatalogsUnderMetalake` now tracks whether it temporarily enabled a
    disabled metalake, and restores the disabled state (best effort) before
    rethrowing on failure, mirroring the enable's own best-effort semantics.
    A vanished metalake (`NoSuchMetalakeException`) still takes the benign
    already-gone path.
    
    ### Why are the changes needed?
    
    A force drop that failed during child cleanup left a user-disabled
    metalake enabled with its catalogs partially dropped, violating the
    disable the user requested.
    
    Fix: #13435
    
    ### Does this PR introduce _any_ user-facing change?
    
    No.
    
    ### How was this patch tested?
    
    Added
    `TestMetalakeManager.testFailedForceDropKeepsDisabledMetalakeDisabled`,
    which pins that a force drop failing mid-cleanup leaves the metalake
    disabled. It fails against the pre-fix code.
---
 .../apache/gravitino/metalake/MetalakeManager.java | 146 ++++++++++++----
 .../gravitino/metalake/TestMetalakeManager.java    | 190 +++++++++++++++++++++
 2 files changed, 299 insertions(+), 37 deletions(-)

diff --git 
a/core/src/main/java/org/apache/gravitino/metalake/MetalakeManager.java 
b/core/src/main/java/org/apache/gravitino/metalake/MetalakeManager.java
index df02f48466..961906bee0 100644
--- a/core/src/main/java/org/apache/gravitino/metalake/MetalakeManager.java
+++ b/core/src/main/java/org/apache/gravitino/metalake/MetalakeManager.java
@@ -29,6 +29,7 @@ import java.util.Arrays;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.stream.Collectors;
 import org.apache.gravitino.Entity.EntityType;
 import org.apache.gravitino.EntityAlreadyExistsException;
@@ -361,37 +362,46 @@ public class MetalakeManager implements 
MetalakeDispatcher, Closeable {
     // Mirror CatalogManager.dropCatalog → ops.dropSchema: force-drop children 
through the real
     // drop path so FilesetCatalogOperations / CatalogManager clean 
write-through secrets (and
     // managed storage). Do this before the metalake root lock to avoid 
nesting tree locks.
+    boolean temporarilyEnabled = false;
     if (force) {
-      dropCatalogsUnderMetalake(ident);
+      temporarilyEnabled = dropCatalogsUnderMetalake(ident);
     }
 
-    return TreeLockUtils.doWithRootTreeLock(
-        LockType.WRITE,
-        () -> {
-          try {
-            boolean inUse = metalakeInUse(store, ident);
-            if (inUse && !force) {
-              throw new MetalakeInUseException(
-                  "Metalake %s is in use, please disable it first or use force 
option", ident);
-            }
+    try {
+      return TreeLockUtils.doWithRootTreeLock(
+          LockType.WRITE,
+          () -> {
+            try {
+              boolean inUse = metalakeInUse(store, ident);
+              if (inUse && !force) {
+                throw new MetalakeInUseException(
+                    "Metalake %s is in use, please disable it first or use 
force option", ident);
+              }
 
-            List<CatalogEntity> catalogEntities =
-                store.list(Namespace.of(ident.name()), CatalogEntity.class, 
EntityType.CATALOG);
-            if (!catalogEntities.isEmpty() && !force) {
-              throw new NonEmptyMetalakeException(
-                  "Metalake %s has catalogs, please drop them first or use 
force option", ident);
-            }
+              List<CatalogEntity> catalogEntities =
+                  store.list(Namespace.of(ident.name()), CatalogEntity.class, 
EntityType.CATALOG);
+              if (!catalogEntities.isEmpty() && !force) {
+                throw new NonEmptyMetalakeException(
+                    "Metalake %s has catalogs, please drop them first or use 
force option", ident);
+              }
 
-            return store.delete(ident, EntityType.METALAKE, true);
-          } catch (NoSuchMetalakeException | NoSuchEntityException e) {
-            // Another server may have completed the drop after the initial 
existence check.
-            // Dropping an already-removed metalake remains an idempotent 
false result.
-            return false;
+              return store.delete(ident, EntityType.METALAKE, true);
+            } catch (NoSuchMetalakeException | NoSuchEntityException e) {
+              // Another server may have completed the drop after the initial 
existence check.
+              // Dropping an already-removed metalake remains an idempotent 
false result.
+              return false;
 
-          } catch (IOException e) {
-            throw new RuntimeException(e);
-          }
-        });
+            } catch (IOException e) {
+              throw new RuntimeException(e);
+            }
+          });
+    } catch (RuntimeException e) {
+      // Phase 1 briefly re-enabled a user-disabled metalake to clean up its 
catalogs. If the
+      // metalake delete above then fails, the metalake still exists, so 
restore its disabled
+      // state rather than leaving the user's disable silently undone.
+      restoreDisabledState(ident, temporarilyEnabled);
+      throw e;
+    }
   }
 
   /**
@@ -401,32 +411,93 @@ public class MetalakeManager implements 
MetalakeDispatcher, Closeable {
    *
    * <p>Callers typically {@code disableMetalake} before force-drop. {@link
    * CatalogManager#dropCatalog} requires catalog {@code 
metalake-in-use=true}, so a disabled
-   * metalake is briefly re-enabled for child cleanup. The metalake entity is 
deleted immediately
-   * afterward, so the temporary enable is not restored.
+   * metalake is briefly re-enabled for child cleanup. If the cleanup fails it 
is re-disabled here
+   * (best effort); if the cleanup succeeds this returns whether the metalake 
was temporarily
+   * enabled, so {@code dropMetalake} can re-disable it should the metalake 
delete then fail. Either
+   * way a user-disabled metalake does not stay enabled after a failed force 
drop.
+   *
+   * @param metalakeIdent the metalake whose child catalogs are force-dropped
+   * @return {@code true} if this temporarily enabled a user-disabled 
metalake, so the caller must
+   *     restore the disabled state if the subsequent metalake delete fails
    */
-  private void dropCatalogsUnderMetalake(NameIdentifier metalakeIdent) {
+  private boolean dropCatalogsUnderMetalake(NameIdentifier metalakeIdent) {
     if (catalogManager == null) {
-      return;
+      return false;
     }
+    // True once this call has flipped the metalake to in-use. Set under the 
metalake write lock
+    // right after the store update (before the non-atomic catalog-status 
update), so a failure in
+    // either the enable itself or the later drops re-disables a metalake this 
operation enabled,
+    // while a concurrent enable leaves it false so we never re-disable one we 
did not enable.
+    AtomicBoolean weEnabled = new AtomicBoolean(false);
     try {
-      if (!metalakeInUse(store, metalakeIdent)) {
-        enableMetalake(metalakeIdent);
-      }
-      List<CatalogEntity> catalogs =
-          store.list(Namespace.of(metalakeIdent.name()), CatalogEntity.class, 
EntityType.CATALOG);
-      for (CatalogEntity catalog : catalogs) {
-        catalogManager.dropCatalog(
-            NameIdentifier.of(metalakeIdent.name(), catalog.name()), true /* 
force */);
+      try {
+        enableMetalakeIfDisabled(metalakeIdent, weEnabled);
+        List<CatalogEntity> catalogs =
+            store.list(Namespace.of(metalakeIdent.name()), 
CatalogEntity.class, EntityType.CATALOG);
+        for (CatalogEntity catalog : catalogs) {
+          catalogManager.dropCatalog(
+              NameIdentifier.of(metalakeIdent.name(), catalog.name()), true /* 
force */);
+        }
+        // Report whether we enabled the metalake so the caller can re-disable 
it if the metalake
+        // delete that follows fails.
+        return weEnabled.get();
+      } catch (NoSuchMetalakeException e) {
+        // Metalake is already gone; dropMetalake will return false. Nothing 
to restore.
+        throw e;
+      } catch (IOException | RuntimeException e) {
+        restoreDisabledState(metalakeIdent, weEnabled.get());
+        throw e;
       }
     } catch (NoSuchMetalakeException e) {
       // Metalake is already gone; dropMetalake will return false.
+      return false;
     } catch (IOException e) {
       throw new RuntimeException(e);
     }
   }
 
+  /**
+   * Best effort: undo a temporary enable this force drop performed, so a 
failed force drop keeps a
+   * user-disabled metalake disabled. Only re-disables when {@code weEnabled} 
is true, i.e. this
+   * operation is the one that enabled the metalake, so it never overwrites a 
concurrent user's
+   * enable.
+   */
+  private void restoreDisabledState(NameIdentifier metalakeIdent, boolean 
weEnabled) {
+    if (!weEnabled) {
+      return;
+    }
+    try {
+      disableMetalake(metalakeIdent);
+    } catch (Exception restoreFailure) {
+      LOG.warn(
+          "Failed to restore the disabled state of metalake {} after a failed 
force drop; "
+              + "the metalake may remain enabled",
+          metalakeIdent,
+          restoreFailure);
+    }
+  }
+
   @Override
   public void enableMetalake(NameIdentifier ident) throws 
NoSuchMetalakeException {
+    enableMetalakeIfDisabled(ident, new AtomicBoolean());
+  }
+
+  /**
+   * Enables the metalake only if it is currently disabled, deciding and 
writing atomically under
+   * the metalake write lock.
+   *
+   * <p>{@code enabledByUs} is set to true the moment this call flips the 
metalake to in-use, before
+   * the non-atomic catalog-status update. A force-drop caller reads it to 
decide whether it must
+   * re-disable the metalake on failure: keying that on this flag (rather than 
a prior {@code
+   * metalakeInUse} read) means a concurrent enable cannot make the caller 
claim, and later undo, an
+   * enable it did not perform, and a failure partway through the enable still 
leaves the caller
+   * able to re-disable.
+   *
+   * @param ident the metalake to enable
+   * @param enabledByUs set to true iff this call performed the disabled to 
in-use flip
+   */
+  private void enableMetalakeIfDisabled(NameIdentifier ident, AtomicBoolean 
enabledByUs)
+      throws NoSuchMetalakeException {
     TreeLockUtils.doWithTreeLock(
         ident,
         LockType.WRITE,
@@ -453,6 +524,7 @@ public class MetalakeManager implements MetalakeDispatcher, 
Closeable {
 
                   return builder.build();
                 });
+            enabledByUs.set(true);
 
             // The only problem is that we can't make sure we can change all 
catalog properties
             // in a transaction. If any catalog property update fails, the 
metalake is already
diff --git 
a/core/src/test/java/org/apache/gravitino/metalake/TestMetalakeManager.java 
b/core/src/test/java/org/apache/gravitino/metalake/TestMetalakeManager.java
index a9a69ec126..d27b096c55 100644
--- a/core/src/test/java/org/apache/gravitino/metalake/TestMetalakeManager.java
+++ b/core/src/test/java/org/apache/gravitino/metalake/TestMetalakeManager.java
@@ -352,6 +352,196 @@ public class TestMetalakeManager {
     store.close();
   }
 
+  @Test
+  public void testFailedForceDropKeepsDisabledMetalakeDisabled() throws 
Exception {
+    CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+    Object originalEnvCatalogManager =
+        FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager", 
true);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager", 
catalogManager, true);
+    try {
+      MetalakeManager manager =
+          new MetalakeManager(entityStore, new RandomIdGenerator(), 
catalogManager);
+      NameIdentifier ident = NameIdentifier.of("metalake_force_drop_failure");
+      manager.createMetalake(ident, "comment", ImmutableMap.of());
+      manager.disableMetalake(ident);
+      Assertions.assertFalse(MetalakeManager.metalakeInUse(entityStore, 
ident));
+
+      entityStore.put(
+          CatalogEntity.builder()
+              .withId(new RandomIdGenerator().nextId())
+              .withName("catalog1")
+              .withNamespace(Namespace.of(ident.name()))
+              .withType(Catalog.Type.RELATIONAL)
+              .withProvider("hive")
+              .withAuditInfo(
+                  
AuditInfo.builder().withCreator("test").withCreateTime(Instant.now()).build())
+              .build(),
+          false);
+
+      Mockito.doThrow(new RuntimeException("catalog drop failed"))
+          .when(catalogManager)
+          .dropCatalog(Mockito.any(), Mockito.anyBoolean());
+
+      Assertions.assertThrows(RuntimeException.class, () -> 
manager.dropMetalake(ident, true));
+
+      // A failed force drop must not leave a user-disabled metalake 
re-enabled.
+      Assertions.assertFalse(
+          MetalakeManager.metalakeInUse(entityStore, ident),
+          "failed force drop must keep the user-disabled metalake disabled");
+    } finally {
+      FieldUtils.writeField(
+          GravitinoEnv.getInstance(), "catalogManager", 
originalEnvCatalogManager, true);
+    }
+  }
+
+  @Test
+  public void testForceDropRestoresDisabledMetalakeWhenDeleteFails() throws 
Exception {
+    CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+    Object originalEnvCatalogManager =
+        FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager", 
true);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager", 
catalogManager, true);
+    try {
+      EntityStore spyStore = Mockito.spy(entityStore);
+      MetalakeManager manager =
+          new MetalakeManager(spyStore, new RandomIdGenerator(), 
catalogManager);
+      NameIdentifier ident = 
NameIdentifier.of("metalake_force_drop_delete_failure");
+      manager.createMetalake(ident, "comment", ImmutableMap.of());
+      manager.disableMetalake(ident);
+      Assertions.assertFalse(MetalakeManager.metalakeInUse(spyStore, ident));
+
+      // No catalogs, so phase-1 catalog cleanup succeeds and re-enables the 
metalake; make the
+      // phase-2 metalake delete fail so the temporary enable outlives the 
cleanup handler.
+      Mockito.doThrow(new IOException("metalake delete failed"))
+          .when(spyStore)
+          .delete(ident, EntityType.METALAKE, true);
+
+      Assertions.assertThrows(RuntimeException.class, () -> 
manager.dropMetalake(ident, true));
+
+      Assertions.assertFalse(
+          MetalakeManager.metalakeInUse(spyStore, ident),
+          "a force drop whose metalake delete fails must keep the 
user-disabled metalake disabled");
+    } finally {
+      FieldUtils.writeField(
+          GravitinoEnv.getInstance(), "catalogManager", 
originalEnvCatalogManager, true);
+    }
+  }
+
+  @Test
+  public void testFailedForceDropDoesNotDisableMetalakeItDidNotEnable() throws 
Exception {
+    CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+    Object originalEnvCatalogManager =
+        FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager", 
true);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager", 
catalogManager, true);
+    try {
+      MetalakeManager manager =
+          new MetalakeManager(entityStore, new RandomIdGenerator(), 
catalogManager);
+      NameIdentifier ident = 
NameIdentifier.of("metalake_force_drop_already_enabled");
+      manager.createMetalake(ident, "comment", ImmutableMap.of());
+      // Left enabled: this stands in for a metalake a concurrent user has 
enabled. The force drop
+      // does not enable it, so a failure must not re-disable it and clobber 
that user's enable.
+      Assertions.assertTrue(MetalakeManager.metalakeInUse(entityStore, ident));
+
+      entityStore.put(
+          CatalogEntity.builder()
+              .withId(new RandomIdGenerator().nextId())
+              .withName("catalog1")
+              .withNamespace(Namespace.of(ident.name()))
+              .withType(Catalog.Type.RELATIONAL)
+              .withProvider("hive")
+              .withAuditInfo(
+                  
AuditInfo.builder().withCreator("test").withCreateTime(Instant.now()).build())
+              .build(),
+          false);
+
+      Mockito.doThrow(new RuntimeException("catalog drop failed"))
+          .when(catalogManager)
+          .dropCatalog(Mockito.any(), Mockito.anyBoolean());
+
+      Assertions.assertThrows(RuntimeException.class, () -> 
manager.dropMetalake(ident, true));
+
+      Assertions.assertTrue(
+          MetalakeManager.metalakeInUse(entityStore, ident),
+          "a failed force drop must not disable a metalake it did not enable");
+    } finally {
+      FieldUtils.writeField(
+          GravitinoEnv.getInstance(), "catalogManager", 
originalEnvCatalogManager, true);
+    }
+  }
+
+  @Test
+  public void testFailedForceDropRestoresWhenTemporaryEnableFailsMidway() 
throws Exception {
+    CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+    Object originalEnvCatalogManager =
+        FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager", 
true);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager", 
catalogManager, true);
+    try {
+      MetalakeManager manager =
+          new MetalakeManager(entityStore, new RandomIdGenerator(), 
catalogManager);
+      NameIdentifier ident = 
NameIdentifier.of("metalake_force_drop_enable_midfail");
+      manager.createMetalake(ident, "comment", ImmutableMap.of());
+      manager.disableMetalake(ident);
+      Assertions.assertFalse(MetalakeManager.metalakeInUse(entityStore, 
ident));
+
+      entityStore.put(
+          CatalogEntity.builder()
+              .withId(new RandomIdGenerator().nextId())
+              .withName("catalog1")
+              .withNamespace(Namespace.of(ident.name()))
+              .withType(Catalog.Type.RELATIONAL)
+              .withProvider("hive")
+              .withAuditInfo(
+                  
AuditInfo.builder().withCreator("test").withCreateTime(Instant.now()).build())
+              .build(),
+          false);
+
+      // The temporary enable flips the metalake in-use, then fails while 
propagating that to the
+      // catalog. A metalake this force drop enabled must not be left enabled.
+      Mockito.doThrow(new RuntimeException("catalog in-use update failed"))
+          .when(catalogManager)
+          .setMetalakeInUseStatus(Mockito.any(), Mockito.anyBoolean());
+
+      Assertions.assertThrows(RuntimeException.class, () -> 
manager.dropMetalake(ident, true));
+
+      Assertions.assertFalse(
+          MetalakeManager.metalakeInUse(entityStore, ident),
+          "a temporary enable that fails partway must not leave the metalake 
enabled");
+    } finally {
+      FieldUtils.writeField(
+          GravitinoEnv.getInstance(), "catalogManager", 
originalEnvCatalogManager, true);
+    }
+  }
+
+  @Test
+  public void testForceDropDeleteFailureDoesNotDisableAlreadyEnabledMetalake() 
throws Exception {
+    CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+    Object originalEnvCatalogManager =
+        FieldUtils.readField(GravitinoEnv.getInstance(), "catalogManager", 
true);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "catalogManager", 
catalogManager, true);
+    try {
+      EntityStore spyStore = Mockito.spy(entityStore);
+      MetalakeManager manager =
+          new MetalakeManager(spyStore, new RandomIdGenerator(), 
catalogManager);
+      NameIdentifier ident = 
NameIdentifier.of("metalake_force_drop_enabled_delete_fail");
+      manager.createMetalake(ident, "comment", ImmutableMap.of());
+      // Left enabled (a concurrent user's enable); no catalogs, so phase 1 
succeeds without this
+      // operation enabling it, then the phase-2 metalake delete fails.
+      Assertions.assertTrue(MetalakeManager.metalakeInUse(spyStore, ident));
+
+      Mockito.doThrow(new IOException("metalake delete failed"))
+          .when(spyStore)
+          .delete(ident, EntityType.METALAKE, true);
+
+      Assertions.assertThrows(RuntimeException.class, () -> 
manager.dropMetalake(ident, true));
+
+      Assertions.assertTrue(
+          MetalakeManager.metalakeInUse(spyStore, ident),
+          "a delete-failed force drop must not disable a metalake it did not 
enable");
+    } finally {
+      FieldUtils.writeField(
+          GravitinoEnv.getInstance(), "catalogManager", 
originalEnvCatalogManager, true);
+    }
+  }
+
   private void testProperties(Map<String, String> expectedProps, Map<String, 
String> testProps) {
     expectedProps.forEach(
         (k, v) -> {

Reply via email to