nsivabalan commented on code in PR #20080:
URL: https://github.com/apache/hudi/pull/20080#discussion_r4171040078


##########
hudi-hadoop-common/src/test/java/org/apache/hudi/hadoop/fs/RecordingLocalFileSystem.java:
##########
@@ -0,0 +1,384 @@
+/*
+ * 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.hudi.hadoop.fs;
+
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.CreateFlag;
+import org.apache.hadoop.fs.FSDataInputStream;
+import org.apache.hadoop.fs.FSDataOutputStream;
+import org.apache.hadoop.fs.FileStatus;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.LocatedFileStatus;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.fs.RawLocalFileSystem;
+import org.apache.hadoop.fs.RemoteIterator;
+import org.apache.hadoop.fs.permission.FsPermission;
+import org.apache.hadoop.util.Progressable;
+
+import java.io.IOException;
+import java.io.PrintWriter;
+import java.io.StringWriter;
+import java.util.Arrays;
+import java.util.EnumSet;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Queue;
+import java.util.WeakHashMap;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.function.BooleanSupplier;
+import java.util.function.Predicate;
+import java.util.stream.Collectors;
+
+/**
+ * The local file system with a recorder of every file system call, for tests 
that assert which
+ * calls a code path makes: how many times it lists the timeline, whether it 
opens a file once per
+ * task or once per job, whether an engine task touches the {@code .hoodie} 
folder at all.
+ *
+ * <p>Each call is recorded with its operation, path and thread, and whether 
it ran in the current
+ * scope, a test-defined notion such as "inside a Spark task" or "inside the 
compaction operator"
+ * (see {@link #setScope}). Calls in scope also keep their call stack. Only 
the outermost call is
+ * recorded, so {@code exists} delegating to {@code getFileStatus} counts 
once. Compound operations
+ * that a store serves as several requests ({@code globStatus}, {@code 
listFiles},
+ * {@code listStatusIterator}) are not recorded themselves; the directory 
listings and status calls
+ * they make are, when they make them. Tests filter the recorded calls with 
the predicates on
+ * {@link Call}, for example {@code 
getCalls(Call.inScope().and(Call.underMetaFolder()))}.
+ *
+ * <p>Usage: {@link #register} the file system on the Hadoop configuration the 
code under test uses
+ * (with Flink options, set {@code "hadoop." + FILE_IMPL_KEY} and {@code 
"hadoop." + DISABLE_CACHE_KEY}
+ * instead), {@link #reset}, run the work and inspect {@link #getCalls}. 
Assert that a call the work
+ * must make was recorded before asserting that another one was not, so that 
the test fails rather
+ * than passes when the recorder is not in use.
+ *
+ * <p>The recorder state is static because Hadoop creates a new instance per 
{@code FileSystem.get}
+ * call once caching is disabled, which is needed for the implementation to be 
picked up regardless
+ * of what the file system cache already holds.
+ */
+public class RecordingLocalFileSystem extends RawLocalFileSystem {
+
+  public static final String FILE_IMPL_KEY = "fs.file.impl";
+  public static final String DISABLE_CACHE_KEY = "fs.file.impl.disable.cache";
+
+  private static final String META_FOLDER = "/.hoodie";
+  private static final int MAX_DESCRIBED_STACKS = 3;
+  private static final BooleanSupplier NO_SCOPE = () -> false;
+
+  private static final Queue<Call> CALLS = new ConcurrentLinkedQueue<>();
+  private static final ThreadLocal<Boolean> IN_RECORDED_CALL = 
ThreadLocal.withInitial(() -> false);
+  private static final Map<Configuration, Map<String, String>> 
VALUES_BEFORE_REGISTER = new WeakHashMap<>();
+  private static volatile BooleanSupplier scope = NO_SCOPE;
+
+  /**
+   * One file system call.
+   */
+  public static final class Call {
+    private final String operation;
+    private final String path;
+    private final String threadName;
+    private final boolean inScope;
+    private final boolean underMetaFolder;
+    private final Throwable stack;
+
+    private Call(String operation, Path path, boolean inScope) {
+      this.operation = operation;
+      this.path = String.valueOf(path);
+      this.threadName = Thread.currentThread().getName();
+      this.inScope = inScope;
+      this.underMetaFolder = isMetaFolderPath(path);
+      this.stack = inScope ? new Throwable(operation + " " + path) : null;
+    }
+
+    public String getOperation() {
+      return operation;
+    }
+
+    public String getPath() {
+      return path;
+    }
+
+    public String getThreadName() {
+      return threadName;
+    }
+
+    public boolean isInScope() {
+      return inScope;
+    }
+
+    public boolean isUnderMetaFolder() {
+      return underMetaFolder;
+    }
+
+    /**
+     * The call stack, kept for calls in scope only.
+     */
+    public Throwable getStack() {
+      return stack;
+    }
+
+    public static Predicate<Call> inScope() {
+      return Call::isInScope;
+    }
+
+    /**
+     * Calls to a table's {@code .hoodie} folder or anything under it, the 
metadata table included.
+     */
+    public static Predicate<Call> underMetaFolder() {
+      return Call::isUnderMetaFolder;
+    }
+
+    public static Predicate<Call> operation(String... operations) {
+      List<String> names = Arrays.asList(operations);
+      return call -> names.contains(call.getOperation());
+    }
+
+    public static Predicate<Call> pathEndsWith(String suffix) {
+      return call -> call.getPath().endsWith(suffix);
+    }
+
+    public static Predicate<Call> pathContains(String part) {
+      return call -> call.getPath().contains(part);
+    }
+
+    @Override
+    public String toString() {
+      return operation + " " + path + " [" + threadName + "]";
+    }
+  }
+
+  /**
+   * Restores the scope that was set before {@link #withScope}.
+   */
+  public static final class Scope implements AutoCloseable {
+    private final BooleanSupplier previous;
+
+    private Scope(BooleanSupplier previous) {
+      this.previous = previous;
+    }
+
+    @Override
+    public void close() {
+      scope = previous;
+    }
+  }
+
+  /**
+   * Makes the {@code file} scheme resolve to this file system on the given 
configuration.
+   */
+  public static void register(Configuration conf) {
+    synchronized (VALUES_BEFORE_REGISTER) {
+      VALUES_BEFORE_REGISTER.computeIfAbsent(conf, c -> {
+        Map<String, String> values = new HashMap<>();
+        values.put(FILE_IMPL_KEY, c.getRaw(FILE_IMPL_KEY));
+        values.put(DISABLE_CACHE_KEY, c.getRaw(DISABLE_CACHE_KEY));
+        return values;
+      });
+    }
+    conf.setClass(FILE_IMPL_KEY, RecordingLocalFileSystem.class, 
FileSystem.class);
+    conf.setBoolean(DISABLE_CACHE_KEY, true);

Review Comment:
   Nit: a one-line comment on why the cache is disabled (so `FileSystem.get` 
hands back this class on the registered conf rather than a cached stock 
instance) would save the next reader a trip.



##########
hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestFileSystemBackedTableMetadata.java:
##########
@@ -229,6 +241,40 @@ public void testStrayFilesAreFilteredOut() throws 
Exception {
         "Stray files should be filtered out by listPartitions");
   }
 
+  /**
+   * The meta folder is never a partition, so listing must neither probe it 
for partition metadata
+   * nor list it as a partition, even when it holds a stray partition metafile.
+   */
+  @Test
+  void testMetaFolderIsNotProbedForPartitionMetadata() throws Exception {
+    hoodieTestTable = hoodieTestTable.addCommit("100");
+    for (String partition : ONE_LEVEL_PARTITIONS) {
+      hoodieTestTable = hoodieTestTable.withPartitionMetaFiles(partition)
+          .withBaseFilesInPartition(partition, IntStream.range(0, 
2).toArray());
+    }
+    try (OutputStream out = metaClient.getStorage().create(
+        new StoragePath(metaClient.getMetaPath(), 
HoodiePartitionMetadata.HOODIE_PARTITION_METAFILE_PREFIX))) {
+      out.write("stray".getBytes());
+    }
+
+    StorageConfiguration<Configuration> conf = 
HadoopFSUtils.getStorageConfWithCopy(
+        metaClient.getStorageConf().unwrapAs(Configuration.class));
+    RecordingLocalFileSystem.register(conf.unwrap());
+    RecordingLocalFileSystem.reset();

Review Comment:
   Non-blocking hygiene: the test never calls `reset()` or `unregister()` after 
itself. Harmless here because the recording file system is only wired into a 
copied configuration, but an `@AfterEach` doing both would make that explicit 
for the next test that copies this pattern.



##########
hudi-hadoop-common/src/test/java/org/apache/hudi/hadoop/fs/RecordingLocalFileSystem.java:
##########
@@ -0,0 +1,384 @@
+/*
+ * 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.hudi.hadoop.fs;
+
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.CreateFlag;
+import org.apache.hadoop.fs.FSDataInputStream;
+import org.apache.hadoop.fs.FSDataOutputStream;
+import org.apache.hadoop.fs.FileStatus;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.LocatedFileStatus;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.fs.RawLocalFileSystem;
+import org.apache.hadoop.fs.RemoteIterator;
+import org.apache.hadoop.fs.permission.FsPermission;
+import org.apache.hadoop.util.Progressable;
+
+import java.io.IOException;
+import java.io.PrintWriter;
+import java.io.StringWriter;
+import java.util.Arrays;
+import java.util.EnumSet;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Queue;
+import java.util.WeakHashMap;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.function.BooleanSupplier;
+import java.util.function.Predicate;
+import java.util.stream.Collectors;
+
+/**
+ * The local file system with a recorder of every file system call, for tests 
that assert which
+ * calls a code path makes: how many times it lists the timeline, whether it 
opens a file once per
+ * task or once per job, whether an engine task touches the {@code .hoodie} 
folder at all.
+ *
+ * <p>Each call is recorded with its operation, path and thread, and whether 
it ran in the current
+ * scope, a test-defined notion such as "inside a Spark task" or "inside the 
compaction operator"
+ * (see {@link #setScope}). Calls in scope also keep their call stack. Only 
the outermost call is
+ * recorded, so {@code exists} delegating to {@code getFileStatus} counts 
once. Compound operations
+ * that a store serves as several requests ({@code globStatus}, {@code 
listFiles},
+ * {@code listStatusIterator}) are not recorded themselves; the directory 
listings and status calls
+ * they make are, when they make them. Tests filter the recorded calls with 
the predicates on
+ * {@link Call}, for example {@code 
getCalls(Call.inScope().and(Call.underMetaFolder()))}.
+ *
+ * <p>Usage: {@link #register} the file system on the Hadoop configuration the 
code under test uses
+ * (with Flink options, set {@code "hadoop." + FILE_IMPL_KEY} and {@code 
"hadoop." + DISABLE_CACHE_KEY}
+ * instead), {@link #reset}, run the work and inspect {@link #getCalls}. 
Assert that a call the work
+ * must make was recorded before asserting that another one was not, so that 
the test fails rather
+ * than passes when the recorder is not in use.
+ *
+ * <p>The recorder state is static because Hadoop creates a new instance per 
{@code FileSystem.get}
+ * call once caching is disabled, which is needed for the implementation to be 
picked up regardless
+ * of what the file system cache already holds.
+ */
+public class RecordingLocalFileSystem extends RawLocalFileSystem {

Review Comment:
   Non-blocking: the call log, scope and registered configs are all static, so 
the helper assumes one test at a time per JVM. That holds today (`forkCount=1`, 
`reuseForks=true`, no JUnit parallel execution in this module), but since 
#20072 / #20087 / #20112 reuse this file, a sentence in the class doc stating 
the single-test-per-JVM assumption would keep a future parallel test from 
silently sharing the log.
   
   Also worth noting in the doc: this extends `RawLocalFileSystem`, not 
`LocalFileSystem`, so registering it as the `file` scheme drops the checksum 
(`.crc`) behavior of the stock local file system. Fine, probably even 
desirable, but the Spark read-path tests in #20112 should know their `file:` 
listings are not the stock ones.



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