Copilot commented on code in PR #13528:
URL: https://github.com/apache/gravitino/pull/13528#discussion_r4115263760


##########
core/src/test/java/org/apache/gravitino/stats/storage/TestLancePartitionStatisticStorage.java:
##########
@@ -56,6 +56,61 @@
 
 public class TestLancePartitionStatisticStorage {
 
+  @Test
+  public void testInFlightReaderSurvivesCacheReplaceEviction() throws 
Exception {
+    // A reader that borrowed the cached dataset must not have it closed 
underneath it when a
+    // concurrent update REPLACES the cache entry: the removal listener marks 
the holder closed,
+    // and the dataset only closes after the last reader leaves.
+    EntityStore entityStore = mock(EntityStore.class);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "entityStore", 
entityStore, true);
+    LancePartitionStatisticStorage storage =
+        new LancePartitionStatisticStorage(java.util.Collections.emptyMap());

Review Comment:
   The new test constructs a `LancePartitionStatisticStorage`, which owns a 
`RootAllocator`, but never closes the storage. This leaks the allocator on 
every run (and also leaks it if an assertion fails); wrap the test body in 
try-with-resources or a finally block that closes `storage`.



##########
core/src/main/java/org/apache/gravitino/stats/storage/LancePartitionStatisticStorage.java:
##########
@@ -642,11 +693,53 @@ void checkoutLatest() {
       dataset.checkoutLatest();
     }
 
+    /**
+     * Registers an in-flight reader that borrowed the dataset; pairs with 
{@link #readerDone()}.
+     */
+    void readerBegins() {
+      activeReaders.incrementAndGet();
+    }
+
+    /** Releases a reader; closes the dataset when the cache already evicted 
it and none remain. */
+    void readerDone() {
+      if (activeReaders.decrementAndGet() == 0 && closed.get()) {
+        closeOnce();
+      }
+    }
+
     @Override
     public void close() throws IOException {
       if (closed.compareAndSet(false, true)) {
+        if (activeReaders.get() == 0) {
+          closeOnce();
+        }
+        // Readers still scanning close it in readerDone() when the last one 
leaves.

Review Comment:
   Deferring `Dataset.close()` is not sufficient during storage shutdown: 
`LancePartitionStatisticStorage.close()` calls `closeDatasetHolder` (which now 
intentionally leaves active holders open) and then closes the shared Arrow 
allocator immediately. An in-flight scanner still uses that allocator, so 
shutdown can invalidate or fail the scan before `readerDone()` gets a chance to 
close the dataset. Coordinate allocator shutdown with the reader count, or 
explicitly disallow/await shutdown while readers are active.



##########
core/src/test/java/org/apache/gravitino/stats/storage/TestLancePartitionStatisticStorage.java:
##########
@@ -56,6 +56,61 @@
 
 public class TestLancePartitionStatisticStorage {
 
+  @Test
+  public void testInFlightReaderSurvivesCacheReplaceEviction() throws 
Exception {
+    // A reader that borrowed the cached dataset must not have it closed 
underneath it when a
+    // concurrent update REPLACES the cache entry: the removal listener marks 
the holder closed,
+    // and the dataset only closes after the last reader leaves.
+    EntityStore entityStore = mock(EntityStore.class);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "entityStore", 
entityStore, true);
+    LancePartitionStatisticStorage storage =
+        new LancePartitionStatisticStorage(java.util.Collections.emptyMap());
+
+    java.util.concurrent.atomic.AtomicBoolean datasetClosed =
+        new java.util.concurrent.atomic.AtomicBoolean(false);
+    LancePartitionStatisticStorage.DatasetHolder holder =
+        new LancePartitionStatisticStorage.DatasetHolder(null) {
+          @Override
+          protected void closeDataset() {
+            datasetClosed.set(true);
+          }
+        };
+    LancePartitionStatisticStorage.DatasetHolder replacement =
+        new LancePartitionStatisticStorage.DatasetHolder(null);
+
+    java.lang.reflect.Field cacheField =
+        LancePartitionStatisticStorage.class.getDeclaredField("datasetCache");
+    cacheField.setAccessible(true);
+    com.github.benmanes.caffeine.cache.Cache<Long, 
LancePartitionStatisticStorage.DatasetHolder>
+        cache =
+            com.github.benmanes.caffeine.cache.Caffeine.newBuilder()
+                .maximumSize(2)
+                .removalListener(
+                    (Long key,
+                        LancePartitionStatisticStorage.DatasetHolder value,
+                        com.github.benmanes.caffeine.cache.RemovalCause cause) 
-> {
+                      if (value != null
+                          && cause != 
com.github.benmanes.caffeine.cache.RemovalCause.EXPLICIT) {
+                        try {
+                          value.close();
+                        } catch (java.io.IOException e) {
+                          throw new RuntimeException(e);
+                        }
+                      }
+                    })
+                .build();
+    cacheField.set(storage, java.util.Optional.of(cache));
+
+    cache.put(1L, holder);
+    holder.readerBegins();
+    cache.put(1L, replacement); // evicts the old entry, listener calls close()

Review Comment:
   The regression is in `withDataset`: its `asMap().compute` must pin the 
holder before a real `listStatistics` scan can begin. This test bypasses that 
path by calling `readerBegins()` directly and then invoking `cache.put`, so it 
would still pass if `withDataset` looked up the holder non-atomically or failed 
to call `readerDone`; please add coverage that blocks an actual scan and 
triggers a replacement/size eviction while it is in flight.



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