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


##########
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:
   [Important] This helper is now the only fenced post-external-drop delete, 
but two other paths that delete a registration after an external drop has 
already succeeded still delete by name, and they are on the Iceberg REST 
service, not a corner case.
   
   `IcebergTableHookDispatcher.dropTable` calls `dispatcher.dropTable(...)` and 
then `bestEffortReconcileTableEntity`, which calls `deleteTableEntity` -> 
`store.delete(ident, Entity.EntityType.TABLE)` with no `expected` 
(`IcebergTableHookDispatcher.java:288-291`). 
`IcebergViewHookDispatcher.dropView` does the same via `deleteViewEntity` 
(`IcebergViewHookDispatcher.java:254-257`). The comment right above each call 
names this PR's race explicitly: "without a distributed TreeLock, another node 
may recreate the same table between the drop above and the EntityStore delete, 
leaving a stale Gravitino entity if we blindly delete" 
(`IcebergTableHookDispatcher.java:105-107`). The mitigation there is a 
`tableExists` probe followed by `importTableEntity` 
(`IcebergTableHookDispatcher.java:263-267`), which is the ABA this PR closes 
elsewhere: if the name was re-created in between, the by-name delete removes 
the *new* registration with its owner, tags and grants, and the re-import then 
registers th
 e object under a fresh id, so the attachments are gone rather than merely 
stale.
   
   This module is already in the change set — the PR updated 
`TestIcebergNamespaceHookDispatcher` for the cleaner's new 4-argument delete — 
and `EntityStore.getVersion` plus the fenced `delete` are available to it. 
Suggest observing the registration before `dispatcher.dropTable(...)` / 
`dispatcher.dropView(...)` and passing it to the fenced delete, the same shape 
as `TableOperationDispatcher.dropTable`. If that is deliberately deferred, it 
is worth saying so in the description, because as it stands the Iceberg REST 
drop path keeps the behaviour the rest of the PR removes.
   
   Verified by: read `IcebergTableHookDispatcher.java:100-300` and 
`IcebergViewHookDispatcher.java:100-268` in full at HEAD `bdc3ca0`; grepped 
every `store.put`/`store.delete` under 
`core/src/main/java/org/apache/gravitino/catalog/` and 
`iceberg/iceberg-rest-server/src/main/java/`; confirmed the two Iceberg call 
sites are unchanged by this PR with `git diff origin/main...HEAD --name-only`, 
which lists only the three Iceberg *test* files.



##########
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:
   [Important] The fence here is id-only, and the import paths can hand a *new* 
object the id this drop observed, so the fence passes and deletes the new 
registration. That is the second hazard the description's "Why" section names 
("Upserting a new object into a stale row can also reuse the old ID"), fixed 
for creates via `putCreatedEntity` but not for imports.
   
   `importTable` still does `store.put(tableEntity, true /* overwrite */)` 
(`TableOperationDispatcher.java:562`), and the overwrite is an upsert on the 
natural key that does **not** assign `table_id`: 
`insertTableMetaOnDuplicateKeyUpdate` updates `table_name`, the parent ids, 
`audit_info`, the versions and `deleted_at`, and leaves `table_id` at the 
stored value (`TableMetaBaseSQLProvider.java:241-251`). The pre-existing test 
`testNaturalKeyOverwriteUsesPersistedTableId` asserts exactly this for 
MySQL/H2: after `insertTable(replacement, true)` with a fresh random id, 
`stored.id()` still equals `original.id()` 
(`TestTableMetaService.java:272-307`).
   
   Failure scenario: store row for `t` has id X. Node A starts `dropTable(t)`, 
observes X, and the external drop succeeds. `t` is then re-created out of band 
in the source catalog. Node B serves `loadTable(t)`, finds no 
`gravitino.identifier`, takes `uid = idGenerator.nextId()` 
(`TableOperationDispatcher.java:541`) and upserts — the row keeps id X. Node A 
resumes: `deleteObservedRegistration` sees id X under the name, 
`checkExpectedIdentity` passes, and the live table's freshly imported 
registration is deleted along with its columns, owner, tags and grants. The 
version no longer guards this, since `OccWriteSupport.checkExpectedIdentity` 
compares ids only and the CAS uses the row's current version.
   
   `importSchema` (`SchemaOperationDispatcher.java:681`), `importTopic` 
(`TopicOperationDispatcher.java:311`) and `importView` 
(`ViewOperationDispatcher.java:578`) have the same shape. Suggest routing the 
import paths through `putCreatedEntity`, or at least not upserting onto a row 
whose id differs from the entity being imported, so an import always produces a 
row whose id identifies the incarnation it describes.
   
   Verified by: read `importTable` at `TableOperationDispatcher.java:513-571` 
and `insertTable` at `TableMetaService.java:111-193` in full at HEAD `bdc3ca0`; 
read the upsert SQL at `TableMetaBaseSQLProvider.java:224-252`; read 
`testNaturalKeyOverwriteUsesPersistedTableId`, which is skipped only for 
PostgreSQL (`TestTableMetaService.java:271-308`); grepped `store.put(` across 
`core/src/main/java/org/apache/gravitino/catalog/` — the four `overwrite = 
true` call sites left are the four import 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:
   [Important] This `throw` contradicts the contract written 40 lines above it, 
and it re-opens the hazard that paragraph is there to prevent.
   
   The javadoc states: "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" (`OperationDispatcher.java:305-307`). Every other exit from the 
loop honours that — a different id, a vanished row, and `observed == null` all 
return `false`. Only the attempt-limit path throws, and 
`OptimisticLockException` is mapped to a 409 by 
`ExceptionHandlers.java:1177-1180` (`Utils.optimisticLockConflict`), which is 
precisely the retryable status the comment argues against.
   
   Failure scenario: `dropTable(t)` observes id X and the external drop 
succeeds, so table `t` is gone from the catalog. Another writer keeps bumping 
X's `current_version` (a tag or owner attachment, a statistic update — anything 
that goes through `updateTable`), so the CAS in `deleteTableWithVersion` loses 
three times. The client gets a 409. It retries `dropTable(t)`; meanwhile 
someone creates a new `t`, and the retry drops that new external table. Note 
the retry does not even buy the intended protection here: after three attempts 
the registration is left behind anyway, so throwing gets the orphaned row *and* 
the retryable error.
   
   Suggest logging at WARN and returning `false` on exhaustion, matching the 
other branches, and leaving the leftover row to the reconcile pass the javadoc 
already defers to. Whichever way you go, the javadoc and the `@throws` tag at 
line 323 should agree with the code.
   
   Verified by: read `OperationDispatcher.java:302-370` at HEAD `bdc3ca0`; 
confirmed the 409 mapping at 
`server/src/main/java/org/apache/gravitino/server/web/rest/ExceptionHandlers.java:1177`;
 grepped `core/src/main/java/org/apache/gravitino/catalog/` for a `catch 
(OptimisticLockException` around the four `deleteObservedRegistration` call 
sites (`TableOperationDispatcher.java:428` and `:477`, 
`SchemaOperationDispatcher.java:597`, `TopicOperationDispatcher.java:258`, 
`ViewOperationDispatcher.java:346`) — none catches it, so it reaches the REST 
layer.



##########
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:
   [Question] Is this probe safe on an asynchronous catalog? Elsewhere in these 
dispatchers the code assumes it is not — `internalCreateTable` skips re-reading 
the table because "some catalog APIs are asynchronous" 
(`TableOperationDispatcher.java:666-667`), and `internalCreateTopic` repeats it 
(`TopicOperationDispatcher.java:376-377`).
   
   If `schemaExists` can still answer `true` shortly after a `dropSchema` that 
returned `true`, this branch treats a perfectly ordinary successful drop as a 
re-creation: the registration and all child registrations are kept, 
`droppedFromStore` stays `false`, and the follow-up 
`SchemaEntityCleaner.deleteOrphanedSchemaEntities` call at line 601 uses the 
same predicate, so its first candidate is `ident` itself, `schemaExists` breaks 
the walk, and it cleans nothing either. The request still returns `true` 
because `droppedFromCatalog` is `true`, so the caller sees success while the 
store keeps the whole stale subtree until some later drop happens to re-probe 
it.
   
   I could not find a catalog in-tree where I can show that happening, which is 
why this is a question rather than a finding. If you have a catalog in mind 
where the probe is reliable, it would help to say so in the comment; if not, 
gating the skip on `droppedFromCatalog == false` (the case where the name 
plausibly belongs to something else) would keep the protection for the 
cascade-over-a-re-created-schema case without making a successful drop depend 
on read-after-write.
   
   Verified by: read `dropSchema` at `SchemaOperationDispatcher.java:525-614` 
and `schemaRecreatedInCatalog` at `:694-703` at HEAD `bdc3ca0`; traced the 
follow-on cleanup through `SchemaEntityCleaner.deleteOrphanedSchemaEntities` 
(`SchemaEntityCleaner.java:58-98`), where `includeSelf = false` still starts 
the walk at the dropped schema's own scope list; compared against the 
async-catalog comments at `TableOperationDispatcher.java:667` and 
`TopicOperationDispatcher.java:377`.



##########
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;
+        }
+        LOG.info(
+            "The {} registration of {} was updated concurrently; retrying the 
delete (attempt {})",
+            type.name().toLowerCase(Locale.ROOT),
+            ident,
+            attempt + 1);
+      } catch (NoSuchEntityException e) {
+        LOG.warn("The {} to be dropped does not exist in the store: {}", type, 
ident, e);
+        return false;
+      } catch (IOException e) {
+        throw new RuntimeException(e);
+      }
+    }
+  }
+
+  /**
+   * Stores the registration of an entity that the external catalog has just 
created.
+   *
+   * <p>A successful external create means a registration observed before the 
create with another id
+   * is stale: its object was dropped out of band, or by a drop on another 
server that has not
+   * reached the store yet. A registration first observed after the create may 
belong to a newer
+   * object, so it must be kept. An upsert by name would keep that row's id 
and hand its owner,
+   * tags, and grants to the new object, and the pending drop's identity fence 
would then pass and
+   * delete the new registration. The stale row is replaced instead: deleted 
fenced on its own id,
+   * then the new registration is inserted.
+   *
+   * @param entity the registration of the newly created object
+   * @param cascade whether a stale registration is deleted with its children
+   * @param observed the registration read before the external create, or null 
if none existed
+   * @param <E> the entity type
+   * @throws IOException if a store operation fails
+   * @throws OptimisticLockException if the stale registration changed while 
it was replaced
+   */
+  protected <E extends Entity & HasIdentifier> void putCreatedEntity(
+      E entity, boolean cascade, @Nullable EntityVersion observed) throws 
IOException {
+    try {
+      store.put(entity, false /* overwrite */);
+      return;
+    } catch (EntityAlreadyExistsException e) {
+      // A registration already sits under the name; decide below whether it 
is stale.
+    }
+
+    NameIdentifier ident = entity.nameIdentifier();
+    EntityVersion existing = observeRegistration(ident, entity.type());
+    if (existing != null && existing.id() == entity.id()) {
+      // Another node already imported this object. Keep any updates it has 
made since then.
+      return;
+    }
+    if (existing != null) {
+      if (observed == null || existing.id() != observed.id()) {
+        throw new OptimisticLockException(
+            "The registration of %s changed during create; keeping the newer 
registration %s",
+            ident, existing);
+      }
+      LOG.warn(
+          "Replacing the stale {} registration {} of {} with the newly created 
object {}",
+          entity.type().name().toLowerCase(Locale.ROOT),
+          existing,
+          ident,
+          entity.id());
+      store.delete(ident, entity.type(), cascade, observed);
+    }
+    store.put(entity, false /* overwrite */);

Review Comment:
   [Nit] The stale-registration replacement is two independent store calls 
where the old code was one atomic upsert, and for schemas the first call 
cascades.
   
   `createSchema` passes `cascade = true` 
(`SchemaOperationDispatcher.java:230`), so line 422 removes the stale schema 
row *and* every child registration under it before line 424 inserts the new 
row. If the insert then fails — an `EntityAlreadyExistsException` from a racer 
that took the name in between, or any `IOException` — the caller only logs and 
returns a combined schema with no entity 
(`SchemaOperationDispatcher.java:231-236`), leaving the just-created schema 
with no registration at all. `loadSchema` will re-import it eventually, so this 
self-heals for the schema row, and the cascaded children were stale by 
construction, which is why this is a nit rather than a finding.
   
   `EntityStore.executeInTransaction` (`EntityStore.java:400`) would make the 
delete and the insert atomic and is the obvious fit here, though I see no other 
dispatcher using it, so it may not be the convention you want in this layer.
   
   Verified by: read `putCreatedEntity` at `OperationDispatcher.java:390-426` 
at HEAD `bdc3ca0`, the four call sites (`SchemaOperationDispatcher.java:230`, 
`TableOperationDispatcher.java:712`, `TopicOperationDispatcher.java:400`, 
`ViewOperationDispatcher.java:433`) and the `catch (Exception e)` handler each 
one wraps it in; confirmed `store.put(entity, false)` surfaces a duplicate name 
as `EntityAlreadyExistsException` through `ExceptionUtils.checkSQLException` in 
`TableMetaService.insertTable` (`TableMetaService.java:189-192`); grepped 
`executeInTransaction` across 
`core/src/main/java/org/apache/gravitino/catalog/` — no current users.



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