yuqi1129 commented on code in PR #13198:
URL: https://github.com/apache/gravitino/pull/13198#discussion_r4123726186


##########
core/src/main/java/org/apache/gravitino/catalog/OperationDispatcher.java:
##########
@@ -244,6 +256,174 @@ protected StringIdentifier 
getStringIdFromProperties(Map<String, String> propert
     }
   }
 
+  /**
+   * Wraps an updater so the store rejects the update, before writing 
anything, when the row under
+   * the name is not the entity the external catalog reported.
+   *
+   * <p>The updater runs inside the store's update, after the current row is 
read and before the
+   * version-checked write. Throwing here therefore aborts the transaction 
with nothing written,
+   * which is what a post-write id comparison cannot do.
+   *
+   * @param expectedId the expected id of the entity being updated
+   * @param updater the update to apply when the ids match
+   * @param <E> the entity type
+   * @return the guarded updater
+   */
+  protected static <E extends Entity & HasIdentifier> Function<E, E> 
requireEntityId(
+      long expectedId, Function<E, E> updater) {
+    return entity -> {
+      if (entity.id() != expectedId) {
+        throw new EntityIdMismatchException(entity.id(), expectedId);
+      }
+      return updater.apply(entity);
+    };
+  }
+
+  /**
+   * Reads the id and store version of a registration before an 
external-catalog call, so a later
+   * store write can be fenced on it.
+   *
+   * @param ident the entity identifier
+   * @param type the entity type
+   * @return the observed id and version, or null when nothing is registered 
under the name
+   * @throws UnsupportedOperationException if the store cannot read versions
+   */
+  @Nullable
+  protected EntityVersion observeRegistration(NameIdentifier ident, 
Entity.EntityType type) {
+    try {
+      return store.getVersion(ident, type);
+    } catch (NoSuchEntityException e) {
+      return null;
+    } catch (IOException e) {
+      throw new RuntimeException("Failed to read the registration of " + 
ident, e);
+    }
+  }
+
+  /**
+   * Deletes the registration observed before an external-catalog call, and 
only that one.
+   *
+   * <p>The external drop has already succeeded when this runs, so a conflict 
is not reported to the
+   * caller: a failed request would invite a retry, and by now the name may 
belong to a newly
+   * created object that the retry would drop. Instead:
+   *
+   * <ul>
+   *   <li>A registration with another id belongs to a newer incarnation. It 
is kept, and a
+   *       reconcile pass, not this drop, decides what happens to it.
+   *   <li>The same id with a newer version means a concurrent store write to 
the entity that was
+   *       just dropped. The fenced delete is retried, because giving up would 
leave a registration
+   *       whose external object is gone, and a retried drop would no longer 
reach the store.
+   * </ul>
+   *
+   * @param ident the entity identifier
+   * @param type the entity type
+   * @param cascade whether to delete the children as well
+   * @param observed the id and version read before the external call, or null 
when there was no
+   *     registration to delete
+   * @return true if the observed registration was deleted
+   * @throws OptimisticLockException if the observed entity kept changing on 
every attempt
+   * @throws UnsupportedOperationException if the store cannot delete with a 
version check
+   */
+  protected boolean deleteObservedRegistration(

Review Comment:
   Fixed in d220e5b436. Iceberg table and view drops now read the stored ID 
before the backend drop. The store delete checks that ID, so it keeps a new 
object created under the same name. Added tests for both paths.



##########
core/src/main/java/org/apache/gravitino/catalog/OperationDispatcher.java:
##########
@@ -244,6 +256,174 @@ protected StringIdentifier 
getStringIdFromProperties(Map<String, String> propert
     }
   }
 
+  /**
+   * Wraps an updater so the store rejects the update, before writing 
anything, when the row under
+   * the name is not the entity the external catalog reported.
+   *
+   * <p>The updater runs inside the store's update, after the current row is 
read and before the
+   * version-checked write. Throwing here therefore aborts the transaction 
with nothing written,
+   * which is what a post-write id comparison cannot do.
+   *
+   * @param expectedId the expected id of the entity being updated
+   * @param updater the update to apply when the ids match
+   * @param <E> the entity type
+   * @return the guarded updater
+   */
+  protected static <E extends Entity & HasIdentifier> Function<E, E> 
requireEntityId(
+      long expectedId, Function<E, E> updater) {
+    return entity -> {
+      if (entity.id() != expectedId) {
+        throw new EntityIdMismatchException(entity.id(), expectedId);
+      }
+      return updater.apply(entity);
+    };
+  }
+
+  /**
+   * Reads the id and store version of a registration before an 
external-catalog call, so a later
+   * store write can be fenced on it.
+   *
+   * @param ident the entity identifier
+   * @param type the entity type
+   * @return the observed id and version, or null when nothing is registered 
under the name
+   * @throws UnsupportedOperationException if the store cannot read versions
+   */
+  @Nullable
+  protected EntityVersion observeRegistration(NameIdentifier ident, 
Entity.EntityType type) {
+    try {
+      return store.getVersion(ident, type);
+    } catch (NoSuchEntityException e) {
+      return null;
+    } catch (IOException e) {
+      throw new RuntimeException("Failed to read the registration of " + 
ident, e);
+    }
+  }
+
+  /**
+   * Deletes the registration observed before an external-catalog call, and 
only that one.
+   *
+   * <p>The external drop has already succeeded when this runs, so a conflict 
is not reported to the
+   * caller: a failed request would invite a retry, and by now the name may 
belong to a newly
+   * created object that the retry would drop. Instead:
+   *
+   * <ul>
+   *   <li>A registration with another id belongs to a newer incarnation. It 
is kept, and a
+   *       reconcile pass, not this drop, decides what happens to it.
+   *   <li>The same id with a newer version means a concurrent store write to 
the entity that was
+   *       just dropped. The fenced delete is retried, because giving up would 
leave a registration
+   *       whose external object is gone, and a retried drop would no longer 
reach the store.
+   * </ul>
+   *
+   * @param ident the entity identifier
+   * @param type the entity type
+   * @param cascade whether to delete the children as well
+   * @param observed the id and version read before the external call, or null 
when there was no
+   *     registration to delete
+   * @return true if the observed registration was deleted
+   * @throws OptimisticLockException if the observed entity kept changing on 
every attempt
+   * @throws UnsupportedOperationException if the store cannot delete with a 
version check
+   */
+  protected boolean deleteObservedRegistration(
+      NameIdentifier ident,
+      Entity.EntityType type,
+      boolean cascade,
+      @Nullable EntityVersion observed) {
+    if (observed == null) {
+      LOG.warn(
+          "No {} registration was found for {} before the external drop; 
leaving the store alone",
+          type.name().toLowerCase(Locale.ROOT),
+          ident);
+      return false;
+    }
+    for (int attempt = 1; ; attempt++) {
+      try {
+        return store.delete(ident, type, cascade, observed);
+      } catch (OptimisticLockException e) {
+        EntityVersion current = observeRegistration(ident, type);
+        if (current == null) {
+          LOG.warn(
+              "The {} registration of {} was removed concurrently while the 
external drop ran",
+              type.name().toLowerCase(Locale.ROOT),
+              ident);
+          return false;
+        }
+        if (current.id() != observed.id()) {
+          LOG.warn(
+              "The {} registration of {} changed while the external drop ran 
(expected {}, found"
+                  + " {}); it is kept instead of being deleted under the new 
incarnation",
+              type.name().toLowerCase(Locale.ROOT),
+              ident,
+              observed,
+              current);
+          return false;
+        }
+        if (attempt >= MAX_FENCED_DELETE_ATTEMPTS) {
+          throw e;

Review Comment:
   Fixed in d220e5b436. After three delete conflicts, we now log a warning and 
keep the stored row instead of throwing an error. The external drop still 
returns success. Updated the Javadoc and added a test for all three attempts 
failing.



##########
core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java:
##########
@@ -424,15 +425,7 @@ public boolean dropTable(NameIdentifier ident) {
           // Gravitino-only metadata. A true out-of-band drop can therefore 
leave a stale
           // registration that requires separate cleanup.
           if (droppedFromCatalog) {
-            try {
-              store.delete(ident, TABLE);
-            } catch (OptimisticLockException e) {
-              throw e;
-            } catch (NoSuchEntityException e) {
-              LOG.warn("The table to be dropped does not exist in the store: 
{}", ident, e);
-            } catch (Exception e) {
-              throw new RuntimeException(e);
-            }
+            deleteObservedRegistration(ident, TABLE, false, observed);

Review Comment:
   Fixed in d220e5b436. All four import paths now use putCreatedEntity instead 
of overwrite. They read the stored ID before loading the external object and 
only replace that old row. A newer row is kept. Added a test for import running 
during a drop.



##########
core/src/main/java/org/apache/gravitino/catalog/SchemaOperationDispatcher.java:
##########
@@ -556,32 +560,41 @@ public boolean dropSchema(NameIdentifier ident, boolean 
cascade) throws NonEmpty
             schemaProperties = new HashMap<>(schemaEntity.properties());
           }
 
+          // For managed schema, we don't need to drop the schema from the 
store again.
+          boolean isManagedSchema = isManagedEntity(catalogIdent, 
Capability.Scope.SCHEMA);
+          // Read the registration before the external call, so the store 
delete below can only
+          // remove the row this drop started with and never one re-created 
under the same name.
+          EntityVersion observed = isManagedSchema ? null : 
observeRegistration(ident, SCHEMA);
           boolean droppedFromCatalog =
               doWithCatalog(
                   catalogIdent,
                   c -> c.doWithSchemaOps(s -> s.dropSchema(ident, cascade)),
                   NonEmptySchemaException.class,
                   RuntimeException.class);
 
-          // For managed schema, we don't need to drop the schema from the 
store again.
-          boolean isManagedSchema = isManagedEntity(catalogIdent, 
Capability.Scope.SCHEMA);
           if (isManagedSchema) {
             if (droppedFromCatalog) {
               secretManager.deleteSecretsFromProperties(schemaProperties);
             }
             return droppedFromCatalog;
           }
 
-          // A non-cascading drop preserves a missing registration because the 
source schema
-          // may have been renamed. An explicit cascading drop also removes 
stale metadata.
+          // A non-cascading false result may mean the external schema was 
renamed, so preserve
+          // its registration. An explicit cascade also removes stale 
metadata, but only if the
+          // registration is still the one observed before the external call.
           boolean droppedFromStore = false;
           if (droppedFromCatalog || cascade) {
-            try {
-              droppedFromStore = store.delete(ident, SCHEMA, true);
-            } catch (NoSuchEntityException e) {
-              LOG.warn("The schema to be dropped does not exist in the store: 
{}", ident, e);
-            } catch (Exception e) {
-              throw new RuntimeException(e);
+            // The cascade removes every child registered under the observed 
row, including one
+            // registered after the observation, which the identity fence on 
the schema cannot
+            // tell apart. A child can only be created while the schema exists 
in the catalog, so
+            // if it exists again the name was re-created meanwhile: keep its 
registrations.
+            if (schemaRecreatedInCatalog(catalogIdent, ident)) {

Review Comment:
   Agreed, true may mean the drop is not visible yet. In d220e5b436, I kept the 
check and made the comment clear. We also keep the schema and its children if 
the check fails. I did not limit the check to failed drops, since a schema can 
be created again after a successful drop too. Added tests for both cases.



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