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


##########
core/src/main/java/org/apache/gravitino/catalog/OperationDispatcher.java:
##########
@@ -244,6 +249,102 @@ 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 id read from the external catalog
+   * @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 or the
+   *     store cannot read versions
+   */
+  @Nullable
+  protected EntityVersion observeRegistration(NameIdentifier ident, 
Entity.EntityType type) {
+    try {
+      return store.getVersion(ident, type);
+    } catch (NoSuchEntityException | UnsupportedOperationException e) {
+      return null;

Review Comment:
   [Important] `observeRegistration` collapses two different situations into 
`null`: "nothing is registered under this name" (`NoSuchEntityException`) and 
"this store cannot read versions" (`UnsupportedOperationException`). 
`deleteObservedRegistration` then returns early on `observed == null` (line 
316), so for a store that does not implement `getVersion` the unconditional 
fallback at line 327 is unreachable — the only way to reach it would be a store 
that implements `getVersion` but not the version-checked `delete`.
   
   That matters because both `EntityStore.getVersion` and 
`RelationalBackend.getVersion` default to throwing 
`UnsupportedOperationException`, and the backend class is instantiated 
reflectively from config (`RelationalEntityStore.java:176-190`, 
`Class.forName(className)`). A `RelationalBackend` that does not override 
`getVersion` would therefore make every external-backed drop — 
`TableOperationDispatcher.java:432` and `:481`, 
`SchemaOperationDispatcher.java:586`, `TopicOperationDispatcher.java:258`, 
`ViewOperationDispatcher.java:348` — stop deleting the store registration 
altogether and only log a warning, instead of falling back to today's 
unconditional delete as the fallback clearly intends. It also silently flips 
`dropTopic`'s return value to `false`, since `droppedFromStore` is what that 
method returns for managed topics.
   
   Suggest distinguishing the two cases, e.g. have `observeRegistration` signal 
"unsupported" separately (a sentinel `EntityVersion`, an `Optional`, or a 
`supportsVersions()` probe) so `deleteObservedRegistration` can take the 
unconditional path for it while still skipping the delete when there is 
genuinely no registration.
   
   Related, on the same helper: the `catch (UnsupportedOperationException)` at 
line 326 wraps the whole `store.delete(ident, type, cascade, expected)` call, 
so a `UnsupportedOperationException` raised anywhere deeper in the store would 
be downgraded to an unconditional delete by name — exactly the write this PR is 
trying to prevent. Narrowing the fallback to a capability check rather than a 
catch would avoid that too.
   
   Verified by: read `OperationDispatcher.java:285-336` at HEAD `cc52211`; read 
the `getVersion` defaults in `EntityStore.java:241-243` and 
`RelationalBackend.java:154-156`; read 
`RelationalEntityStore.createRelationalEntityBackend` at lines 176-190; and 
grepped `implements EntityStore|implements RelationalBackend` across the repo 
(only `RelationalEntityStore` and `JDBCBackend` in main, plus two test stores).



##########
core/src/main/java/org/apache/gravitino/catalog/SchemaOperationDispatcher.java:
##########
@@ -556,33 +559,31 @@ 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);
-            }
+            droppedFromStore = deleteObservedRegistration(ident, SCHEMA, true, 
observed);

Review Comment:
   [Question] The cascade here is fenced on the schema row only. 
`SchemaMetaService.deleteSchema(identifier, cascade, expected)` compares 
`expected` against the schema PO's `(schemaId, currentVersion)` and then, on 
`cascade`, deletes the children unconditionally.
   
   Inserting a table into a schema does not bump that schema's 
`current_version`, so a table registered after `observeRegistration` at line 
566 and before this call passes the fence and is deleted, even though no drop 
ever observed it. The schema re-creation case is genuinely covered (a new 
schema gets a new id, so the check fires), so this is the same ABA one level 
down rather than a regression — is it in scope for this PR, or deliberately 
left to the reconcile pass the helper's javadoc mentions?
   
   Verified by: read `SchemaMetaService.deleteSchema` at 
`SchemaMetaService.java:302-320` and the cascade branch that follows, plus 
`SchemaOperationDispatcher.java:562-587`, at HEAD `cc52211`.



##########
core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java:
##########
@@ -311,6 +312,25 @@ public boolean delete(NameIdentifier ident, 
Entity.EntityType entityType, boolea
     }
   }
 
+  @Override
+  public EntityVersion getVersion(NameIdentifier ident, Entity.EntityType 
entityType)
+      throws IOException {
+    // Always read through: a cached entity does not carry the store version, 
and a stale cache
+    // entry must never be the basis of a version check.
+    return backend.getVersion(ident, entityType);
+  }
+
+  @Override
+  public boolean delete(
+      NameIdentifier ident, Entity.EntityType entityType, boolean cascade, 
EntityVersion expected)
+      throws IOException {
+    try {
+      return backend.delete(ident, entityType, cascade, expected);
+    } finally {
+      cache.invalidate(ident, entityType);

Review Comment:
   [Important] This is the only mutating path in `RelationalEntityStore` that 
calls `cache.invalidate(...)` directly instead of the `invalidateCache(...)` 
helper, so it skips `cacheInvalidationEpoch.incrementAndGet()`.
   
   That epoch exists precisely to close this window: `batchGet` samples the 
epoch before reading the backend (line 273) and refuses to write a fetched 
entity back into the cache if the epoch moved in between (lines 279, 288-294). 
Because the version-checked delete never advances the epoch, a `batchGet` that 
started before it can repopulate the cache with the just-deleted entity after 
`cache.invalidate` ran, and the stale copy then survives until the TTL. Since 
this new method is exactly the one the drop paths now use 
(`OperationDispatcher.deleteObservedRegistration` -> `store.delete(ident, type, 
cascade, expected)`), the practical effect is that a successfully dropped 
table/schema/topic/view can stay readable from the cache — which undoes part of 
what the PR is fixing.
   
   Every other mutating method here uses the helper: lines 165, 222, 233, 311 
(the unconditional `delete`), 344 (`deleteAndGet`), 452-453, 473-475, 501, 534, 
557.
   
   Suggested fix: `invalidateCache(ident, entityType);`
   
   Verified by: read `RelationalEntityStore.java` at HEAD `cc52211` — the 
`invalidateCache` definition at lines 603-607, the epoch checks in `batchGet` 
at 273-294, and grepped every `invalidateCache|cache.invalidate` occurrence in 
the file; line 330 is the only direct `cache.invalidate` call outside the 
helper.



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