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


##########
core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java:
##########
@@ -223,9 +251,21 @@ public <E extends Entity & HasIdentifier> List<E> batchGet(
                   return entity.isEmpty();
                 })
             .toList();
+    // Unlike get(), the backend read is not done under the entries' cache 
locks: holding one lock
+    // per key across a batch DB round trip would stall unrelated reads on the 
same segments. So an
+    // invalidation can land between the read and the write-back. The epoch 
sampled here detects
+    // that and skips the write-back, otherwise the stale copy would survive 
until the TTL. The
+    // per-key lock makes the check and the put atomic against an invalidation 
of the same key.

Review Comment:
   [Question] "The per-key lock makes the check and the put atomic against an 
invalidation of the same key" holds for `invalidate()` — 
`CaffeineEntityCache.invalidate` (CaffeineEntityCache.java:168) takes the same 
stripe — but it does not extend to `clear()`. `SegmentedLock.withGlobalLock` 
(SegmentedLock.java:201) never acquires the segment locks: it sets a latch and 
synchronizes on itself, and `waitForGlobalComplete()` is only consulted on lock 
*entry* (SegmentedLock.java:85/111/139/165). A write-back that has already 
passed the epoch check is therefore not excluded from 
`cacheData.invalidateAll()`, so it can still land afterwards and leave exactly 
the stale entry the fence targets.
   
   The window is nanoseconds and is the same class as the 
hierarchical-invalidation window the description already calls out — but the 
description lists the clear fallback as fully fenced, so it seems worth naming 
clear() in that caveat (and perhaps in this comment) rather than letting a 
later reader assume the global lock serializes it. Is that the intent, or did 
you read `withGlobalLock` as taking the stripes?
   
   Verified by: read `SegmentedLock.withGlobalLock` / `waitForGlobalComplete` 
and `CaffeineEntityCache.clear` / `invalidate` at this commit.



##########
core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java:
##########
@@ -223,9 +251,21 @@ public <E extends Entity & HasIdentifier> List<E> batchGet(
                   return entity.isEmpty();
                 })
             .toList();
+    // Unlike get(), the backend read is not done under the entries' cache 
locks: holding one lock
+    // per key across a batch DB round trip would stall unrelated reads on the 
same segments. So an
+    // invalidation can land between the read and the write-back. The epoch 
sampled here detects
+    // that and skips the write-back, otherwise the stale copy would survive 
until the TTL. The
+    // per-key lock makes the check and the put atomic against an invalidation 
of the same key.
+    long epochBeforeRead = cacheInvalidationEpoch.get();
     List<E> fetchEntities = backend.batchGet(noCacheIdents, entityType);
     for (E entity : fetchEntities) {
-      cache.put(entity);
+      cache.withCacheLock(
+          EntityCacheKey.of(entity.nameIdentifier(), entity.type()),
+          () -> {
+            if (cacheInvalidationEpoch.get() == epochBeforeRead) {
+              cache.put(entity);
+            }
+          });

Review Comment:
   [Nit] The write-back loop now takes a segment lock per fetched entity even 
when nothing can be cached. `BaseEntityCache.put` returns early for every type 
outside `CACHEABLE_TYPES` (BaseEntityCache.java:116-118; the set at 
BaseEntityCache.java:55-66 contains no USER/GROUP/ROLE), and the one caller 
that is not a cache preload fetches ROLE entities on the authorization path 
(`server-common/.../jcasbin/JcasbinAuthorizer.java:1386`). For that caller each 
entity now costs a `waitForGlobalComplete()` plus a stripe acquire/release 
around a no-op `put`.
   
   Cheap fix: hoist the cacheability test above the loop, the way the other 
call site already guards before calling (`MetadataAuthzHelper.java:615`), e.g. 
skip the `withCacheLock` wrapper when 
`!BaseEntityCache.isCacheable(entityType)`. Keeps the fence where it matters 
and leaves the hot path as it was.
   
   Verified by: read `BaseEntityCache.put` and `CACHEABLE_TYPES`, then grepped 
every `batchGet(` call site in this checkout — `MetadataAuthzHelper` guards on 
`isCacheable`, `JcasbinAuthorizer` does not.



##########
core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java:
##########
@@ -75,6 +76,11 @@ public class RelationalEntityStore
   private EntityChangeLogCleaner entityChangeLogCleaner;
   private EntityCache cache;
 
+  // Advanced before every cache invalidation and clear, whether local or 
replayed from the change
+  // log, so that batchGet() can tell that an invalidation happened while its 
backend read was in
+  // flight. Every invalidation must go through invalidateCache() / 
clearCache() for this to hold.
+  private final AtomicLong cacheInvalidationEpoch = new AtomicLong();

Review Comment:
   [Question] The epoch is a per-process counter, so it only observes 
invalidations this JVM performs. For a `Coherence.LOCAL_PER_NODE` cache that is 
complete, because the change-log listener replays remote invalidations through 
the store (line 136). For a `Coherence.SHARED` cache no listener is registered 
at all (`registerCacheChangeLogListener`, lines 128-137), so another node 
invalidating a shared entry never advances this node's epoch, and `batchGet` 
here would happily write its stale copy back into the shared cache.
   
   No `SHARED` implementation exists in tree today, but `CacheFactory` loads 
the cache class from config, so it is a supported extension point. Would you 
mind noting on this field that the fence is LOCAL_PER_NODE-only (or NONE), so 
whoever adds a shared cache does not inherit a silently inert guard?
   
   Verified by: grepped for `Coherence.SHARED` in `core/src/main` (only the 
enum, Javadoc and the listener-gating comments), and read 
`CacheFactory.getEntityCache`'s reflective class loading.



##########
core/src/test/java/org/apache/gravitino/storage/relational/TestRelationalEntityStoreBatchGetLateFill.java:
##########
@@ -0,0 +1,146 @@
+/*
+ * 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.gravitino.storage.relational;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+
+import java.util.List;
+import org.apache.commons.lang3.reflect.FieldUtils;
+import org.apache.gravitino.Config;
+import org.apache.gravitino.Entity;
+import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.Namespace;
+import org.apache.gravitino.cache.CaffeineEntityCache;
+import org.apache.gravitino.meta.TableEntity;
+import org.apache.gravitino.storage.relational.po.cache.EntityChangeRecord;
+import org.apache.gravitino.storage.relational.po.cache.OperateType;
+import org.apache.gravitino.utils.TestUtil;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+/**
+ * Verifies that {@link RelationalEntityStore#batchGet} cannot write an entity 
back into the cache
+ * after that entity was invalidated while the backend read was in flight (the 
late-fill race).
+ *
+ * <p>Each race is modelled deterministically by running the invalidation 
inside the backend stub:
+ * the backend has already produced its (now stale) result, and the 
invalidation lands before {@code
+ * batchGet} writes it back.
+ */
+public class TestRelationalEntityStoreBatchGetLateFill {
+
+  private static final Namespace SCHEMA_NS = Namespace.of("metalake", 
"catalog", "schema");
+
+  private RelationalEntityStore store;
+  private RelationalBackend backend;
+  private CaffeineEntityCache cache;
+
+  @BeforeEach
+  void setUp() throws IllegalAccessException {
+    store = new RelationalEntityStore();
+    backend = Mockito.mock(RelationalBackend.class);
+    cache = Mockito.spy(new CaffeineEntityCache(new Config() {}));
+    FieldUtils.writeField(store, "backend", backend, true);
+    FieldUtils.writeField(store, "cache", cache, true);
+  }
+
+  private static EntityChangeRecord dropRecord(NameIdentifier ident, 
Entity.EntityType type) {
+    return new EntityChangeRecord(
+        1L,
+        ident.namespace().level(0),
+        type.name(),
+        EntityChangeLogNameIdentifierCodec.encode(ident),
+        OperateType.DROP,
+        0L);
+  }
+
+  @Test
+  void testBatchGetWritesBackWhenNoInvalidationHappens() {
+    TableEntity table = TestUtil.getTestTableEntity(1L, "t1", SCHEMA_NS);
+    Mockito.when(backend.batchGet(any(), 
eq(Entity.EntityType.TABLE))).thenReturn(List.of(table));
+
+    List<TableEntity> result =
+        store.batchGet(List.of(table.nameIdentifier()), 
Entity.EntityType.TABLE, TableEntity.class);
+
+    Assertions.assertEquals(List.of(table), result);
+    Assertions.assertTrue(cache.contains(table.nameIdentifier(), 
Entity.EntityType.TABLE));
+  }
+
+  @Test
+  void testBatchGetSkipsWriteBackWhenChangeLogInvalidatesDuringBackendRead() {
+    TableEntity table = TestUtil.getTestTableEntity(1L, "t1", SCHEMA_NS);
+    NameIdentifier ident = table.nameIdentifier();
+    EntityChangeLogListener poller = store.newCacheChangeLogListener();
+    Mockito.when(backend.batchGet(any(), eq(Entity.EntityType.TABLE)))
+        .thenAnswer(
+            invocation -> {
+              poller.onEntityChange(List.of(dropRecord(ident, 
Entity.EntityType.TABLE)));
+              return List.of(table);
+            });
+
+    List<TableEntity> result =
+        store.batchGet(List.of(ident), Entity.EntityType.TABLE, 
TableEntity.class);
+
+    Assertions.assertEquals(List.of(table), result);
+    Assertions.assertFalse(
+        cache.contains(ident, Entity.EntityType.TABLE),
+        "a value invalidated during the backend read must not be written 
back");
+  }
+
+  @Test
+  void testBatchGetSkipsWriteBackWhenLocalDeleteInvalidatesDuringBackendRead() 
{
+    TableEntity table = TestUtil.getTestTableEntity(1L, "t1", SCHEMA_NS);
+    NameIdentifier ident = table.nameIdentifier();
+    Mockito.when(backend.batchGet(any(), eq(Entity.EntityType.TABLE)))
+        .thenAnswer(
+            invocation -> {
+              store.delete(ident, Entity.EntityType.TABLE, false);
+              return List.of(table);
+            });
+
+    store.batchGet(List.of(ident), Entity.EntityType.TABLE, TableEntity.class);
+
+    Assertions.assertFalse(cache.contains(ident, Entity.EntityType.TABLE));
+  }
+
+  @Test
+  void 
testBatchGetSkipsWriteBackWhenListenerFallsBackToClearDuringBackendRead() {
+    TableEntity table = TestUtil.getTestTableEntity(1L, "t1", SCHEMA_NS);
+    NameIdentifier ident = table.nameIdentifier();
+    EntityChangeLogListener poller = store.newCacheChangeLogListener();
+    // A failed targeted invalidation makes the listener clear the whole 
cache; that must fence
+    // in-flight fills too.
+    Mockito.doThrow(new RuntimeException("boom"))
+        .when(cache)
+        .invalidate(ident, Entity.EntityType.TABLE);
+    Mockito.when(backend.batchGet(any(), eq(Entity.EntityType.TABLE)))
+        .thenAnswer(
+            invocation -> {
+              poller.onEntityChange(List.of(dropRecord(ident, 
Entity.EntityType.TABLE)));
+              return List.of(table);
+            });
+
+    store.batchGet(List.of(ident), Entity.EntityType.TABLE, TableEntity.class);
+
+    Mockito.verify(cache).clear();
+    Assertions.assertFalse(cache.contains(ident, Entity.EntityType.TABLE));

Review Comment:
   [Nit] This test does not isolate the case it names. `invalidateCache` 
advances the epoch *before* calling `cache.invalidate` 
(RelationalEntityStore.java:554-556), so by the time the stubbed 
`RuntimeException` propagates the fence has already been tripped; the clear 
fallback that follows bumps it a second time. Delete the `incrementAndGet()` 
from `clearCache()` (RelationalEntityStore.java:560) and both assertions here 
still pass — that increment is the one line of the change with no test that can 
fail on it.
   
   Driving the clear path on its own would fix it, e.g. make `clearCache()` 
`@VisibleForTesting` package-private and call it from inside the backend stub 
instead of going through a failing invalidation.
   
   Verified by: traced the call order listener -> `Target.invalidate` -> 
`invalidateCache` (epoch++, then the stubbed throw) -> catch -> `Target.clear` 
-> `clearCache` (epoch++), against the epoch check at 
RelationalEntityStore.java:265.



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