This is an automated email from the ASF dual-hosted git repository.

shunping pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 94972054878 Route the remaining GCS calls through GcsUtilV2 when 
enabled (#40273)
94972054878 is described below

commit 949720548789dbedd751fab7eb8f0be5221ec233
Author: Shunping Huang <[email protected]>
AuthorDate: Mon Sep 28 12:22:49 2026 -0400

    Route the remaining GCS calls through GcsUtilV2 when enabled (#40273)
    
    * Route the remaining GCS calls through GcsUtilV2
    
    The facade only had V2 branches on its GcsPath-typed methods, but
    GcsFileSystem and GcpOptions call the String-typed and legacy-model ones,
    so match, list, copy, rename, delete and bucket creation all stayed on V1
    even with use_gcsutil_v2 set. Add V2 branches to those, converting Blob
    and BucketInfo back to the JSON API model so their callers need no change.
    
    Strategies are picked to match V1: copy and rename rewrite without a
    destination precondition, remove ignores a 404, and buckets are created
    with projectPrivate ACLs.
    
    * Match GcsUtilV1's upload chunk size in GcsUtilV2
    
    The formula is copied rather than read from AsyncWriteChannelOptions
    because GcsUtilV2 has no other gcsio reference and the migration ends
    with that dependency deleted. The test keeps the two values in sync.
    
    * Fix a bug of use_gcsutil_v2 flag not passing to experiments in ParquetIOLT
    
    * Fix GcsUtilV2 parity gaps.
    
    - Translate StorageException thrown on write channel close.
    - Record a 404 request metric when open() hits a missing object.
    
    * Fix bucketAccessible in V2 to match V1
    
    * Route the deprecated getBucket and create(path, type) through GcsUtilV2
    
    * Add some more routing tests and trim unnecessary ones.
    
    * Make GcsUtilV2's storage client mockable in tests
    
    Route all client access through a @VisibleForTesting storage() accessor.
    In GcsUtilTest, merge the V2 helpers into one, test bucketAccessible
    through a mocked client, and reset metrics containers in tearDown.
    
    * Add reference link for default upload chunk size
---
 .../apache/beam/it/gcp/storage/ParquetIOLT.java    |  47 +-
 .../beam/sdk/extensions/gcp/util/GcsUtil.java      | 218 +++++++-
 .../beam/sdk/extensions/gcp/util/GcsUtilV2.java    |  94 +++-
 .../beam/sdk/extensions/gcp/util/GcsUtilTest.java  | 559 +++++++++++++++++++++
 .../sdk/extensions/gcp/util/GcsUtilV2Test.java     |  12 +
 5 files changed, 885 insertions(+), 45 deletions(-)

diff --git 
a/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/ParquetIOLT.java
 
b/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/ParquetIOLT.java
index 80eb7886c95..52b12372dcd 100644
--- 
a/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/ParquetIOLT.java
+++ 
b/it/google-cloud-platform/src/test/java/org/apache/beam/it/gcp/storage/ParquetIOLT.java
@@ -31,7 +31,9 @@ import java.nio.ByteBuffer;
 import java.time.Duration;
 import java.time.ZoneOffset;
 import java.time.format.DateTimeFormatter;
+import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.List;
 import java.util.Map;
 import java.util.Random;
 import java.util.UUID;
@@ -120,9 +122,9 @@ import org.junit.runners.MethodSorters;
  *
  * <p>Whichever is chosen, the job is launched with an explicit experiment, 
{@code use_runner_v2} or
  * {@code disable_runner_v2}. Leaving the choice to the service is not an 
option here because the
- * container image is resolved on the client, see {@link 
#dataflowWorkerExperiment()}: a job
- * submitted without an experiment ends up asking for an image tag that does 
not exist and hangs
- * with the workers in ImagePullBackOff.
+ * container image is resolved on the client, see {@link 
#dataflowExperiments()}: a job submitted
+ * without an experiment ends up asking for an image tag that does not exist 
and hangs with the
+ * workers in ImagePullBackOff.
  *
  * <p>Note that Runner v2 stages the locally built SDK jars, so a local SDK 
change is measured as
  * is, while the legacy worker runs the Beam code baked into its container 
image.
@@ -198,7 +200,7 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
   /**
    * Experiment that keeps the job on the legacy worker. Needed even though 
the legacy worker is
    * what a job without experiments is submitted as, because the service 
upgrades such a job to
-   * Runner v2 on its own, see {@link #dataflowWorkerExperiment()}.
+   * Runner v2 on its own, see {@link #dataflowExperiments()}.
    */
   private static final String LEGACY_WORKER_EXPERIMENT = "disable_runner_v2";
 
@@ -499,8 +501,8 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
       // the gcs_* metrics measure.
       // maxNumWorkers is deliberately not set, it only bounds an autoscaling 
pool.
       builder
-          // Picks the worker, see dataflowWorkerExperiment().
-          .addParameter("experiments", dataflowWorkerExperiment())
+          // Picks the worker and the GcsUtil version, see 
dataflowExperiments().
+          .addParameter("experiments", dataflowExperiments())
           .addParameter("autoscalingAlgorithm", "NONE")
           .addParameter("numWorkers", String.valueOf(DATAFLOW_NUM_WORKERS))
           .addParameter("workerMachineType", DATAFLOW_MACHINE_TYPE);
@@ -510,14 +512,19 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
   }
 
   /**
-   * Experiment that selects the Dataflow worker, {@code use_runner_v2} or 
{@code
-   * disable_runner_v2}.
+   * Every experiment the Dataflow job is submitted with, as the comma 
separated list the {@code
+   * experiments} pipeline option parses.
    *
-   * <p>The worker is always selected explicitly, even though Runner v2 is 
what the service picks on
-   * its own, because the container image is resolved on the client: {@code
-   * DataflowRunner.getDefaultContainerImageUrl} takes the image name and the 
image tag from the
-   * same branch of its {@code useUnifiedWorker()} check, and the two tags 
({@code
-   * dataflowFnapiContainerVersion} and {@code dataflowLegacyContainerVersion} 
in {@code
+   * <p>They have to be joined into one value because the launcher carries the 
parameters in a map,
+   * so {@code experiments} can only be given once. That map is also the only 
channel that reaches
+   * the job: the launcher rebuilds the pipeline options from these parameters 
alone, so an
+   * experiment set on the options of the pipeline would be dropped.
+   *
+   * <p>The list always names the worker, {@code use_runner_v2} or {@code 
disable_runner_v2}, even
+   * though Runner v2 is what the service picks on its own, because the 
container image is resolved
+   * on the client: {@code DataflowRunner.getDefaultContainerImageUrl} takes 
the image name and the
+   * image tag from the same branch of its {@code useUnifiedWorker()} check, 
and the two tags
+   * ({@code dataflowFnapiContainerVersion} and {@code 
dataflowLegacyContainerVersion} in {@code
    * runners/google-cloud-dataflow-java/build.gradle}) are bumped 
independently.
    *
    * <p>A job submitted without an experiment therefore resolves the legacy 
pair {@code
@@ -529,8 +536,13 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
    * gives {@code beam_javaNN_sdk:<fnapi tag>}, {@code disable_runner_v2} 
keeps the service from
    * upgrading the job so {@code beam-javaNN-batch:<legacy tag>} stays correct.
    */
-  private static String dataflowWorkerExperiment() {
-    return configuration.useRunnerV2 ? RUNNER_V2_EXPERIMENT : 
LEGACY_WORKER_EXPERIMENT;
+  private static String dataflowExperiments() {
+    List<String> experiments = new ArrayList<>();
+    experiments.add(configuration.useRunnerV2 ? RUNNER_V2_EXPERIMENT : 
LEGACY_WORKER_EXPERIMENT);
+    if (configuration.useGcsUtilV2) {
+      experiments.add(GCS_UTIL_V2_EXPERIMENT);
+    }
+    return String.join(",", experiments);
   }
 
   /**
@@ -669,7 +681,7 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
         DATAFLOW_RUNNER.equalsIgnoreCase(configuration.runner)
             ? String.format(
                 "%s (--experiments=%s)",
-                configuration.useRunnerV2 ? "Runner v2" : "legacy", 
dataflowWorkerExperiment())
+                configuration.useRunnerV2 ? "Runner v2" : "legacy", 
dataflowExperiments())
             : "n/a",
         DATAFLOW_RUNNER.equalsIgnoreCase(configuration.runner)
             ? String.format("%d x %s, autoscaling off", DATAFLOW_NUM_WORKERS, 
DATAFLOW_MACHINE_TYPE)
@@ -875,8 +887,7 @@ public final class ParquetIOLT extends GcsIOLoadTestBase {
     /**
      * Dataflow only. {@code true} runs the job on Runner v2, i.e. the unified 
worker, {@code false}
      * on the legacy worker. Either way the choice is sent to the service as 
an explicit experiment,
-     * see {@link ParquetIOLT#dataflowWorkerExperiment()} for why it must not 
be left to the
-     * service.
+     * see {@link ParquetIOLT#dataflowExperiments()} for why it must not be 
left to the service.
      */
     @JsonProperty public boolean useRunnerV2 = true;
 
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java
index 9e2d82ae1eb..382ac6856bf 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtil.java
@@ -17,6 +17,7 @@
  */
 package org.apache.beam.sdk.extensions.gcp.util;
 
+import com.google.api.client.util.DateTime;
 import com.google.api.gax.paging.Page;
 import com.google.api.services.storage.model.Bucket;
 import com.google.api.services.storage.model.Objects;
@@ -28,10 +29,17 @@ import com.google.cloud.storage.Storage.BlobListOption;
 import com.google.cloud.storage.Storage.BlobSourceOption;
 import com.google.cloud.storage.Storage.BlobWriteOption;
 import com.google.cloud.storage.Storage.BucketGetOption;
+import com.google.cloud.storage.Storage.BucketTargetOption;
+import com.google.cloud.storage.Storage.PredefinedAcl;
+import com.google.cloud.storage.StorageClass;
 import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
 import java.io.IOException;
+import java.math.BigInteger;
 import java.nio.channels.SeekableByteChannel;
 import java.nio.channels.WritableByteChannel;
+import java.time.Duration;
+import java.time.OffsetDateTime;
+import java.util.ArrayList;
 import java.util.Collection;
 import java.util.List;
 import java.util.Set;
@@ -151,6 +159,9 @@ public class GcsUtil {
    */
   @Deprecated
   public StorageObject getObject(GcsPath gcsPath) throws IOException {
+    if (delegateV2 != null) {
+      return toStorageObject(delegateV2.getBlob(gcsPath));
+    }
     return delegate.getObject(gcsPath);
   }
 
@@ -166,6 +177,21 @@ public class GcsUtil {
    */
   @Deprecated
   public List<StorageObjectOrIOException> getObjects(List<GcsPath> gcsPaths) 
throws IOException {
+    if (delegateV2 != null) {
+      List<StorageObjectOrIOException> results = new ArrayList<>();
+      for (BlobResult blobResult : delegateV2.getBlobs(gcsPaths)) {
+        Blob blob = blobResult.blob();
+        IOException ioException = blobResult.ioException();
+        if (blob != null) {
+          
results.add(StorageObjectOrIOException.create(toStorageObject(blob)));
+        } else if (ioException != null) {
+          results.add(StorageObjectOrIOException.create(ioException));
+        } else {
+          throw new IOException("Invalid blob result: it holds neither a blob 
nor an error.");
+        }
+      }
+      return results;
+    }
     List<GcsUtilV1.StorageObjectOrIOException> legacy = 
delegate.getObjects(gcsPaths);
     return legacy.stream()
         .map(StorageObjectOrIOException::fromLegacy)
@@ -186,6 +212,9 @@ public class GcsUtil {
   @Deprecated
   public Objects listObjects(String bucket, String prefix, @Nullable String 
pageToken)
       throws IOException {
+    if (delegateV2 != null) {
+      return toObjects(delegateV2.listBlobs(bucket, prefix, pageToken));
+    }
     return delegate.listObjects(bucket, prefix, pageToken);
   }
 
@@ -196,6 +225,9 @@ public class GcsUtil {
   public Objects listObjects(
       String bucket, String prefix, @Nullable String pageToken, @Nullable 
String delimiter)
       throws IOException {
+    if (delegateV2 != null) {
+      return toObjects(delegateV2.listBlobs(bucket, prefix, pageToken, 
delimiter));
+    }
     return delegate.listObjects(bucket, prefix, pageToken, delimiter);
   }
 
@@ -240,7 +272,9 @@ public class GcsUtil {
    */
   @Deprecated
   public WritableByteChannel create(GcsPath path, String type) throws 
IOException {
-    return delegate.create(path, type);
+    // Built the same way as GcsUtilV1#create(GcsPath, String), but through 
this class so that it
+    // follows GcsUtilV2 when enabled.
+    return create(path, CreateOptions.builder().setContentType(type).build());
   }
 
   /**
@@ -249,7 +283,12 @@ public class GcsUtil {
   @Deprecated
   public WritableByteChannel create(GcsPath path, String type, Integer 
uploadBufferSizeBytes)
       throws IOException {
-    return delegate.create(path, type, uploadBufferSizeBytes);
+    return create(
+        path,
+        CreateOptions.builder()
+            .setContentType(type)
+            .setUploadBufferSizeBytes(uploadBufferSizeBytes)
+            .build());
   }
 
   public static class CreateOptions {
@@ -341,16 +380,27 @@ public class GcsUtil {
   }
 
   /**
-   * @deprecated use {@link #createBucket(BucketInfo)}.
+   * @deprecated use {@link #createBucket(BucketInfo, BucketTargetOption...)}.
    */
   @Deprecated
   public void createBucket(String projectId, Bucket bucket) throws IOException 
{
+    if (delegateV2 != null) {
+      // GcsUtilV1 always creates buckets with projectPrivate ACLs, which 
java-storage does not do
+      // on its own, so they have to be requested explicitly to keep the same 
access.
+      delegateV2.createBucket(
+          projectId,
+          toBucketInfo(bucket),
+          BucketTargetOption.predefinedAcl(PredefinedAcl.PROJECT_PRIVATE),
+          
BucketTargetOption.predefinedDefaultObjectAcl(PredefinedAcl.PROJECT_PRIVATE));
+      return;
+    }
     delegate.createBucket(projectId, bucket);
   }
 
-  public void createBucket(BucketInfo bucketInfo) throws IOException {
+  public void createBucket(BucketInfo bucketInfo, BucketTargetOption... 
options)
+      throws IOException {
     if (delegateV2 != null) {
-      delegateV2.createBucket(bucketInfo);
+      delegateV2.createBucket(bucketInfo, options);
     } else {
       throw new IOException("GcsUtil V2 not initialized.");
     }
@@ -361,6 +411,9 @@ public class GcsUtil {
    */
   @Deprecated
   public @Nullable Bucket getBucket(GcsPath path) throws IOException {
+    if (delegateV2 != null) {
+      return toBucket(delegateV2.getBucket(path));
+    }
     return delegate.getBucket(path);
   }
 
@@ -377,6 +430,10 @@ public class GcsUtil {
    */
   @Deprecated
   public void removeBucket(Bucket bucket) throws IOException {
+    if (delegateV2 != null) {
+      delegateV2.removeBucket(toBucketInfo(bucket));
+      return;
+    }
     delegate.removeBucket(bucket);
   }
 
@@ -390,6 +447,14 @@ public class GcsUtil {
 
   public void copy(Iterable<String> srcFilenames, Iterable<String> 
destFilenames)
       throws IOException {
+    if (delegateV2 != null) {
+      // GcsUtilV1 issues a rewrite without any destination precondition, so 
ALWAYS_OVERWRITE is
+      // the strategy that preserves its behavior. The strategies that inspect 
the destination
+      // would also cost an extra GET per file.
+      delegateV2.copy(
+          toGcsPaths(srcFilenames), toGcsPaths(destFilenames), 
OverwriteStrategy.ALWAYS_OVERWRITE);
+      return;
+    }
     delegate.copy(srcFilenames, destFilenames);
   }
 
@@ -412,6 +477,22 @@ public class GcsUtil {
   public void rename(
       Iterable<String> srcFilenames, Iterable<String> destFilenames, 
MoveOptions... moveOptions)
       throws IOException {
+    GcsUtilV2 v2 = delegateV2;
+    if (v2 != null) {
+      Set<MoveOptions> moveOptionSet = Sets.newHashSet(moveOptions);
+      // Note this differs from renameV2, which defaults to SAFE_OVERWRITE. 
GcsUtilV1 rewrites
+      // without a destination precondition, so ALWAYS_OVERWRITE is the 
behavior preserving choice.
+      v2.move(
+          toGcsPaths(srcFilenames),
+          toGcsPaths(destFilenames),
+          moveOptionSet.contains(StandardMoveOptions.IGNORE_MISSING_FILES)
+              ? MissingStrategy.SKIP_IF_MISSING
+              : MissingStrategy.FAIL_IF_MISSING,
+          
moveOptionSet.contains(StandardMoveOptions.SKIP_IF_DESTINATION_EXISTS)
+              ? OverwriteStrategy.SKIP_IF_EXISTS
+              : OverwriteStrategy.ALWAYS_OVERWRITE);
+      return;
+    }
     delegate.rename(srcFilenames, destFilenames, moveOptions);
   }
 
@@ -453,6 +534,11 @@ public class GcsUtil {
   }
 
   public void remove(Collection<String> filenames) throws IOException {
+    if (delegateV2 != null) {
+      // GcsUtilV1 ignores a 404 on delete, which is SKIP_IF_MISSING.
+      delegateV2.remove(toGcsPaths(filenames), 
MissingStrategy.SKIP_IF_MISSING);
+      return;
+    }
     delegate.remove(filenames);
   }
 
@@ -470,6 +556,128 @@ public class GcsUtil {
     }
   }
 
+  private static List<GcsPath> toGcsPaths(Iterable<String> filenames) {
+    List<GcsPath> paths = new ArrayList<>();
+    for (String filename : filenames) {
+      paths.add(GcsPath.fromUri(filename));
+    }
+    return paths;
+  }
+
+  /**
+   * Converts a JSON API {@link Bucket} into the java-storage {@link 
BucketInfo} model.
+   *
+   * <p>Only the properties that callers of the deprecated {@link 
#createBucket(String, Bucket)} set
+   * are carried over. Like {@link #toStorageObject}, this is expected to go 
away with the
+   * deprecated methods it serves.
+   */
+  private static BucketInfo toBucketInfo(Bucket bucket) {
+    BucketInfo.Builder builder = BucketInfo.newBuilder(bucket.getName());
+    if (bucket.getLocation() != null) {
+      builder.setLocation(bucket.getLocation());
+    }
+    if (bucket.getStorageClass() != null) {
+      builder.setStorageClass(StorageClass.valueOf(bucket.getStorageClass()));
+    }
+    Bucket.SoftDeletePolicy softDeletePolicy = bucket.getSoftDeletePolicy();
+    if (softDeletePolicy != null && 
softDeletePolicy.getRetentionDurationSeconds() != null) {
+      builder.setSoftDeletePolicy(
+          BucketInfo.SoftDeletePolicy.newBuilder()
+              .setRetentionDuration(
+                  
Duration.ofSeconds(softDeletePolicy.getRetentionDurationSeconds()))
+              .build());
+    }
+    return builder.build();
+  }
+
+  /**
+   * Converts a java-storage {@link com.google.cloud.storage.Bucket} back into 
the JSON API {@link
+   * Bucket} model, for the deprecated {@link #getBucket(GcsPath)}.
+   *
+   * <p>The reverse of {@link #toBucketInfo}, plus the owning project number. 
Like it, this is
+   * expected to go away with the deprecated methods it serves.
+   */
+  private static Bucket toBucket(com.google.cloud.storage.Bucket bucketInfo) {
+    Bucket bucket =
+        new Bucket()
+            .setName(bucketInfo.getName())
+            .setLocation(bucketInfo.getLocation())
+            .setProjectNumber(bucketInfo.getProject());
+    StorageClass storageClass = bucketInfo.getStorageClass();
+    if (storageClass != null) {
+      bucket.setStorageClass(storageClass.name());
+    }
+    BucketInfo.SoftDeletePolicy softDeletePolicy = 
bucketInfo.getSoftDeletePolicy();
+    Duration retention = softDeletePolicy == null ? null : 
softDeletePolicy.getRetentionDuration();
+    if (retention != null) {
+      bucket.setSoftDeletePolicy(
+          new 
Bucket.SoftDeletePolicy().setRetentionDurationSeconds(retention.getSeconds()));
+    }
+    return bucket;
+  }
+
+  /**
+   * Converts a java-storage {@link Blob} back into the JSON API {@link 
StorageObject} model.
+   *
+   * <p>This lets the deprecated, legacy typed methods of this class be served 
by {@link GcsUtilV2}
+   * without their callers having to change. It is expected to go away once 
those methods do.
+   */
+  private static StorageObject toStorageObject(Blob blob) {
+    StorageObject storageObject =
+        new StorageObject()
+            .setBucket(blob.getBucket())
+            .setName(blob.getName())
+            .setGeneration(blob.getGeneration())
+            .setMetageneration(blob.getMetageneration())
+            .setContentType(blob.getContentType())
+            .setContentEncoding(blob.getContentEncoding())
+            .setMd5Hash(blob.getMd5())
+            .setCrc32c(blob.getCrc32c())
+            .setEtag(blob.getEtag());
+    Long size = blob.getSize();
+    if (size != null) {
+      storageObject.setSize(BigInteger.valueOf(size));
+    }
+    OffsetDateTime updated = blob.getUpdateTimeOffsetDateTime();
+    if (updated != null) {
+      storageObject.setUpdated(new 
DateTime(updated.toInstant().toEpochMilli()));
+    }
+    OffsetDateTime created = blob.getCreateTimeOffsetDateTime();
+    if (created != null) {
+      storageObject.setTimeCreated(new 
DateTime(created.toInstant().toEpochMilli()));
+    }
+    return storageObject;
+  }
+
+  /** Converts a single page of java-storage {@link Blob}s into the JSON API 
{@link Objects}. */
+  private static Objects toObjects(Page<Blob> page) {
+    List<StorageObject> items = new ArrayList<>();
+    List<String> prefixes = new ArrayList<>();
+    for (Blob blob : page.getValues()) {
+      // A delimited listing reports each common prefix as a directory 
placeholder.
+      if (blob.isDirectory()) {
+        prefixes.add(blob.getName());
+      } else {
+        items.add(toStorageObject(blob));
+      }
+    }
+    Objects objects = new Objects();
+    // Leave items and prefixes null when empty, as the JSON API does, so that 
callers looping on
+    // getItems() != null keep terminating.
+    if (!items.isEmpty()) {
+      objects.setItems(items);
+    }
+    if (!prefixes.isEmpty()) {
+      objects.setPrefixes(prefixes);
+    }
+    // Page.getNextPageToken() may be an empty string rather than null on the 
last page.
+    String nextPageToken = page.hasNextPage() ? page.getNextPageToken() : null;
+    if (nextPageToken != null) {
+      objects.setNextPageToken(nextPageToken);
+    }
+    return objects;
+  }
+
   @SuppressFBWarnings("NM_CLASS_NOT_EXCEPTION")
   public static class StorageObjectOrIOException {
     final GcsUtilV1.StorageObjectOrIOException delegate;
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
index 9d223f81050..20260aa3c79 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2.java
@@ -45,6 +45,7 @@ import com.google.cloud.storage.Storage.BlobSourceOption;
 import com.google.cloud.storage.Storage.BlobWriteOption;
 import com.google.cloud.storage.Storage.BucketField;
 import com.google.cloud.storage.Storage.BucketGetOption;
+import com.google.cloud.storage.Storage.BucketTargetOption;
 import com.google.cloud.storage.Storage.CopyRequest;
 import com.google.cloud.storage.StorageBatch;
 import com.google.cloud.storage.StorageBatchResult;
@@ -100,7 +101,16 @@ class GcsUtilV2 {
     }
   }
 
-  private Storage storage;
+  private final Storage storage;
+
+  /**
+   * The shared client. Every operation reaches it through this method, so 
that tests can substitute
+   * a mocked client for a spied instance.
+   */
+  @VisibleForTesting
+  Storage storage() {
+    return storage;
+  }
 
   private final @Nullable Integer uploadBufferSizeBytes;
 
@@ -116,6 +126,15 @@ class GcsUtilV2 {
   /** Maximum number of requests permitted in a GCS batch request. */
   private static final int MAX_REQUESTS_PER_BATCH = 100;
 
+  /**
+   * Upload chunk size applied when the pipeline does not ask for one. Mirrors 
gcsio's {@code
+   * AsyncWriteChannelOptions} default, which java-storage does not share. 
Reference link:
+   * 
https://github.com/GoogleCloudDataproc/hadoop-connectors/blob/v3.1.14/util/src/main/java/com/google/cloud/hadoop/util/AsyncWriteChannelOptions.java#L72
+   */
+  @VisibleForTesting
+  static final int DEFAULT_UPLOAD_CHUNK_SIZE_BYTES =
+      Runtime.getRuntime().maxMemory() < 512 * 1024 * 1024 ? 8 * 1024 * 1024 : 
3 * 8 * 1024 * 1024;
+
   /**
    * Limit the number of bytes Cloud Storage will attempt to copy before 
responding to an individual
    * request. If you see Read Timeout errors, try reducing this value.
@@ -294,13 +313,13 @@ class GcsUtilV2 {
   @VisibleForTesting
   Storage storageWithHttpMetrics(@Nullable MetricsContainer container, boolean 
isWrite) {
     if (container == null) {
-      return storage;
+      return storage();
     }
-    StorageOptions options = storage.getOptions();
+    StorageOptions options = storage().getOptions();
     TransportOptions transportOptions = options.getTransportOptions();
     if (!(transportOptions instanceof HttpTransportOptions)) {
       // A non-HTTP transport (e.g. gRPC) has no HttpRequestInitializer to 
wrap.
-      return storage;
+      return storage();
     }
     return options.toBuilder()
         .setTransportOptions(
@@ -346,7 +365,7 @@ class GcsUtilV2 {
   }
 
   public Blob getBlob(GcsPath gcsPath, BlobGetOption... options) throws 
IOException {
-    return getBlob(storage, gcsPath, options);
+    return getBlob(storage(), gcsPath, options);
   }
 
   /** As {@link #getBlob(GcsPath, BlobGetOption...)}, but issued through a 
specific client. */
@@ -399,7 +418,7 @@ class GcsUtilV2 {
         Lists.partition(Lists.newArrayList(gcsPaths), MAX_REQUESTS_PER_BATCH)) 
{
 
       // Create a new empty batch every time
-      StorageBatch batch = storage.batch();
+      StorageBatch batch = storage().batch();
       List<StorageBatchResult<Blob>> batchResultFutures = new ArrayList<>();
 
       for (GcsPath path : pathPartition) {
@@ -455,7 +474,7 @@ class GcsUtilV2 {
     }
 
     try {
-      return storage.list(bucket, blobListOptions.toArray(new 
BlobListOption[0]));
+      return storage().list(bucket, blobListOptions.toArray(new 
BlobListOption[0]));
     } catch (StorageException e) {
       throw translateStorageException(bucket, prefix, e);
     }
@@ -525,7 +544,7 @@ class GcsUtilV2 {
         Lists.partition(Lists.newArrayList(paths), MAX_REQUESTS_PER_BATCH)) {
 
       // Create a new empty batch every time
-      StorageBatch batch = storage.batch();
+      StorageBatch batch = storage().batch();
       List<StorageBatchResult<Boolean>> batchResultFutures = new ArrayList<>();
 
       for (GcsPath path : pathPartition) {
@@ -592,7 +611,7 @@ class GcsUtilV2 {
         // FAIL_IF_EXISTS, SKIP_IF_EXISTS and SAFE_OVERWRITE require checking 
the target blob
         BlobInfo existingTarget;
         try {
-          existingTarget = storage.get(dstId);
+          existingTarget = storage().get(dstId);
         } catch (StorageException e) {
           throw translateStorageException(dstPath, e);
         }
@@ -622,11 +641,11 @@ class GcsUtilV2 {
       }
 
       try {
-        CopyWriter copyWriter = storage.copy(copyRequestBuilder.build());
+        CopyWriter copyWriter = storage().copy(copyRequestBuilder.build());
         copyWriter.getResult();
 
         if (deleteSrc) {
-          if (!storage.delete(srcId)) {
+          if (!storage().delete(srcId)) {
             // This may happen if the source file is deleted by another 
process after copy.
             LOG.warn(
                 "Source file {} could not be deleted after move to {}. It may 
not have existed.",
@@ -664,7 +683,7 @@ class GcsUtilV2 {
   public Bucket getBucket(GcsPath path, BucketGetOption... options) throws 
IOException {
     String bucketName = path.getBucket();
     try {
-      Bucket bucket = storage.get(bucketName, options);
+      Bucket bucket = storage().get(bucketName, options);
       if (bucket == null) {
         throw new FileNotFoundException(
             String.format("The specified bucket does not exist: gs://%s", 
bucketName));
@@ -675,13 +694,17 @@ class GcsUtilV2 {
     }
   }
 
-  /** Returns whether the GCS bucket exists and is accessible. */
-  public boolean bucketAccessible(GcsPath path) {
+  /**
+   * Returns whether the GCS bucket exists and is accessible. This will return 
false if the bucket
+   * does not exist or is inaccessible due to permissions; any other failure 
is propagated, as
+   * {@link GcsUtilV1} does.
+   */
+  public boolean bucketAccessible(GcsPath path) throws IOException {
     try {
       // Fetch only the name field to minimize data transfer
       getBucket(path, BucketGetOption.fields(BucketField.NAME));
       return true;
-    } catch (IOException e) {
+    } catch (AccessDeniedException | FileNotFoundException e) {
       return false;
     }
   }
@@ -705,9 +728,27 @@ class GcsUtilV2 {
     return bucket.getProject().longValue();
   }
 
-  public void createBucket(BucketInfo bucketInfo) throws IOException {
+  public void createBucket(BucketInfo bucketInfo, BucketTargetOption... 
options)
+      throws IOException {
+    createBucket(null, bucketInfo, options);
+  }
+
+  /**
+   * As {@link #createBucket(BucketInfo, BucketTargetOption...)}, but creates 
the bucket in the
+   * given project instead of the one this instance is configured with.
+   */
+  public void createBucket(
+      @Nullable String projectId, BucketInfo bucketInfo, BucketTargetOption... 
options)
+      throws IOException {
+    Storage client = storage();
+    if (projectId != null && !projectId.equals(this.projectId)) {
+      // The owning project is a property of the client rather than of the 
insert request, so
+      // asking for a different one means deriving a client for it. The 
derived client shares the
+      // credentials, host and transport of the original.
+      client = 
storage().getOptions().toBuilder().setProjectId(projectId).build().getService();
+    }
     try {
-      storage.create(bucketInfo);
+      client.create(bucketInfo, options);
     } catch (StorageException e) {
       throw translateStorageException(bucketInfo.getName(), null, e);
     }
@@ -715,7 +756,7 @@ class GcsUtilV2 {
 
   public void removeBucket(BucketInfo bucketInfo) throws IOException {
     try {
-      if (!storage.delete(bucketInfo.getName())) {
+      if (!storage().delete(bucketInfo.getName())) {
         throw new FileNotFoundException(
             String.format("The specified bucket does not exist: gs://%s", 
bucketInfo.getName()));
       }
@@ -801,6 +842,11 @@ class GcsUtilV2 {
       serviceCallMetric.call("ok");
       return wrapInCounting(
           new GcsSeekableByteChannel(reader, blob.getSize()), 
path.getBucket(), container);
+    } catch (FileNotFoundException e) {
+      // getBlob reports a missing object as a FileNotFoundException rather 
than a
+      // StorageException, so record its status here like GcsUtilV1 does.
+      serviceCallMetric.call(404);
+      throw e;
     } catch (StorageException e) {
       serviceCallMetric.call(e.getCode());
       throw translateStorageException(path, e);
@@ -833,7 +879,12 @@ class GcsUtilV2 {
 
     @Override
     public void close() throws IOException {
-      writer.close();
+      // The upload is finalized here, so this is where a failed precondition 
surfaces.
+      try {
+        writer.close();
+      } catch (StorageException e) {
+        throw translateStorageException(gcsPath, e);
+      }
     }
   }
 
@@ -873,9 +924,8 @@ class GcsUtilV2 {
           options.getUploadBufferSizeBytes() != null
               ? options.getUploadBufferSizeBytes()
               : this.uploadBufferSizeBytes;
-      if (uploadBufferSizeBytes != null) {
-        writer.setChunkSize(uploadBufferSizeBytes);
-      }
+      writer.setChunkSize(
+          uploadBufferSizeBytes != null ? uploadBufferSizeBytes : 
DEFAULT_UPLOAD_CHUNK_SIZE_BYTES);
 
       serviceCallMetric.call("ok");
       // Return the bridge wrapper
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilTest.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilTest.java
index 52bb877bb96..4beb56d6c20 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilTest.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilTest.java
@@ -55,18 +55,26 @@ import com.google.api.client.testing.http.MockHttpTransport;
 import com.google.api.client.testing.http.MockLowLevelHttpRequest;
 import com.google.api.client.testing.http.MockLowLevelHttpResponse;
 import com.google.api.client.util.BackOff;
+import com.google.api.gax.paging.Page;
 import com.google.api.services.storage.Storage;
 import com.google.api.services.storage.model.Bucket;
 import com.google.api.services.storage.model.Objects;
 import com.google.api.services.storage.model.RewriteResponse;
 import com.google.api.services.storage.model.StorageObject;
 import com.google.auth.Credentials;
+import com.google.cloud.WriteChannel;
 import com.google.cloud.hadoop.gcsio.CreateObjectOptions;
 import com.google.cloud.hadoop.gcsio.GoogleCloudStorage;
 import com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl;
 import com.google.cloud.hadoop.gcsio.GoogleCloudStorageOptions;
 import com.google.cloud.hadoop.gcsio.GoogleCloudStorageReadOptions;
 import com.google.cloud.hadoop.gcsio.StorageResourceId;
+import com.google.cloud.storage.BucketInfo;
+import com.google.cloud.storage.Storage.BucketGetOption;
+import com.google.cloud.storage.Storage.BucketTargetOption;
+import com.google.cloud.storage.Storage.PredefinedAcl;
+import com.google.cloud.storage.StorageClass;
+import com.google.cloud.storage.StorageException;
 import java.io.ByteArrayInputStream;
 import java.io.ByteArrayOutputStream;
 import java.io.FileNotFoundException;
@@ -114,9 +122,11 @@ import org.apache.beam.sdk.util.FluentBackoff;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Lists;
 import org.checkerframework.checker.nullness.qual.Nullable;
+import org.junit.After;
 import org.junit.Before;
 import org.junit.Rule;
 import org.junit.Test;
+import org.junit.function.ThrowingRunnable;
 import org.junit.rules.ExpectedException;
 import org.junit.runner.RunWith;
 import org.junit.runners.JUnit4;
@@ -136,6 +146,13 @@ public class GcsUtilTest {
     MetricsEnvironment.setCurrentContainer(testMetricsContainer);
   }
 
+  @After
+  public void tearDown() {
+    // Don't leak the containers installed by setUp into later tests in the 
same JVM.
+    MetricsEnvironment.setProcessWideContainer(null);
+    MetricsEnvironment.setCurrentContainer(null);
+  }
+
   private static GcsOptions gcsOptionsWithTestCredential() {
     GcsOptions pipelineOptions = PipelineOptionsFactory.as(GcsOptions.class);
     pipelineOptions.setGcpCredential(new TestCredential());
@@ -1887,4 +1904,546 @@ public class GcsUtilTest {
   private static InputStream toStream(String content) throws IOException {
     return new ByteArrayInputStream(content.getBytes(StandardCharsets.UTF_8));
   }
+
+  // The tests below cover the routing that the GcsUtil facade performs for 
the legacy typed
+  // methods once the use_gcsutil_v2 experiment installs a GcsUtilV2 delegate. 
Both delegates are
+  // mocked, so they assert both that V2 is used and that V1 is left alone.
+
+  private GcsUtilV1 mockDelegate;
+  private GcsUtilV2 mockDelegateV2;
+
+  private GcsUtil gcsUtilRoutingToV2() {
+    GcsUtil gcsUtil = gcsOptionsWithTestCredential().getGcsUtil();
+    mockDelegate = Mockito.mock(GcsUtilV1.class);
+    mockDelegateV2 = Mockito.mock(GcsUtilV2.class);
+    gcsUtil.delegate = mockDelegate;
+    gcsUtil.delegateV2 = mockDelegateV2;
+    return gcsUtil;
+  }
+
+  private static com.google.cloud.storage.Blob mockBlob(String bucket, String 
object) {
+    com.google.cloud.storage.Blob blob = 
Mockito.mock(com.google.cloud.storage.Blob.class);
+    when(blob.getBucket()).thenReturn(bucket);
+    when(blob.getName()).thenReturn(object);
+    return blob;
+  }
+
+  @Test
+  public void testCopyIsRoutedToV2AsAnUnconditionalOverwrite() throws 
IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+    gcsUtil.copy(ImmutableList.of("gs://bucket/from"), 
ImmutableList.of("gs://bucket/to"));
+
+    verify(mockDelegateV2)
+        .copy(
+            ImmutableList.of(GcsPath.fromUri("gs://bucket/from")),
+            ImmutableList.of(GcsPath.fromUri("gs://bucket/to")),
+            GcsUtilV2.OverwriteStrategy.ALWAYS_OVERWRITE);
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testRenameWithoutOptionsIsRoutedToV2AsAnUnconditionalOverwrite() 
throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+    gcsUtil.rename(ImmutableList.of("gs://bucket/from"), 
ImmutableList.of("gs://bucket/to"));
+
+    verify(mockDelegateV2)
+        .move(
+            ImmutableList.of(GcsPath.fromUri("gs://bucket/from")),
+            ImmutableList.of(GcsPath.fromUri("gs://bucket/to")),
+            GcsUtilV2.MissingStrategy.FAIL_IF_MISSING,
+            GcsUtilV2.OverwriteStrategy.ALWAYS_OVERWRITE);
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testRenameMoveOptionsAreTranslatedForV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+    gcsUtil.rename(
+        ImmutableList.of("gs://bucket/from"),
+        ImmutableList.of("gs://bucket/to"),
+        StandardMoveOptions.IGNORE_MISSING_FILES,
+        StandardMoveOptions.SKIP_IF_DESTINATION_EXISTS);
+
+    verify(mockDelegateV2)
+        .move(
+            ImmutableList.of(GcsPath.fromUri("gs://bucket/from")),
+            ImmutableList.of(GcsPath.fromUri("gs://bucket/to")),
+            GcsUtilV2.MissingStrategy.SKIP_IF_MISSING,
+            GcsUtilV2.OverwriteStrategy.SKIP_IF_EXISTS);
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testRemoveIsRoutedToV2AndToleratesMissingFiles() throws 
IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+    gcsUtil.remove(ImmutableList.of("gs://bucket/one", "gs://bucket/two"));
+
+    verify(mockDelegateV2)
+        .remove(
+            ImmutableList.of(
+                GcsPath.fromUri("gs://bucket/one"), 
GcsPath.fromUri("gs://bucket/two")),
+            GcsUtilV2.MissingStrategy.SKIP_IF_MISSING);
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  /**
+   * Blobs and errors are passed through in order. The field-by-field 
conversion is covered by
+   * {@link #testGetObjectIsRoutedToV2AndKeepsAllFields}.
+   */
+  @Test
+  public void testGetObjectsIsRoutedToV2AndConvertsBlobs() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    com.google.cloud.storage.Blob blob = mockBlob("bucket", "found");
+    FileNotFoundException notFound = new 
FileNotFoundException("gs://bucket/missing");
+    List<GcsPath> paths =
+        ImmutableList.of(GcsPath.fromUri("gs://bucket/found"), 
GcsPath.fromUri("gs://bucket/miss"));
+    when(mockDelegateV2.getBlobs(paths))
+        .thenReturn(
+            ImmutableList.of(
+                GcsUtilV2.BlobResult.create(blob), 
GcsUtilV2.BlobResult.create(notFound)));
+
+    List<StorageObjectOrIOException> results = gcsUtil.getObjects(paths);
+
+    assertEquals(2, results.size());
+    StorageObject converted = results.get(0).storageObject();
+    assertNotNull(converted);
+    assertEquals("found", converted.getName());
+    assertNull(results.get(0).ioException());
+    assertSame(notFound, results.get(1).ioException());
+    assertNull(results.get(1).storageObject());
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testListObjectsIsRoutedToV2AndConvertsAPage() throws IOException 
{
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    com.google.cloud.storage.Blob object = mockBlob("bucket", "prefix/object");
+    com.google.cloud.storage.Blob directory = mockBlob("bucket", 
"prefix/dir/");
+    when(directory.isDirectory()).thenReturn(true);
+    @SuppressWarnings("unchecked")
+    Page<com.google.cloud.storage.Blob> page = Mockito.mock(Page.class);
+    when(page.getValues()).thenReturn(ImmutableList.of(object, directory));
+    when(page.hasNextPage()).thenReturn(true);
+    when(page.getNextPageToken()).thenReturn("next");
+    when(mockDelegateV2.listBlobs("bucket", "prefix/", null)).thenReturn(page);
+
+    Objects objects = gcsUtil.listObjects("bucket", "prefix/", null);
+
+    assertEquals(1, objects.getItems().size());
+    assertEquals("prefix/object", objects.getItems().get(0).getName());
+    assertEquals(ImmutableList.of("prefix/dir/"), objects.getPrefixes());
+    assertEquals("next", objects.getNextPageToken());
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testListObjectsReportsTheLastPageWithANullToken() throws 
IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    @SuppressWarnings("unchecked")
+    Page<com.google.cloud.storage.Blob> page = Mockito.mock(Page.class);
+    when(page.getValues()).thenReturn(ImmutableList.of());
+    // A gax page reports an empty token rather than a null one once it is 
exhausted. Callers of
+    // listObjects loop until the token is null, so it has to be normalized.
+    when(page.hasNextPage()).thenReturn(false);
+    when(page.getNextPageToken()).thenReturn("");
+    when(mockDelegateV2.listBlobs("bucket", "prefix/", null)).thenReturn(page);
+
+    Objects objects = gcsUtil.listObjects("bucket", "prefix/", null);
+
+    assertNull(objects.getItems());
+    assertNull(objects.getPrefixes());
+    assertNull(objects.getNextPageToken());
+  }
+
+  /**
+   * The storage class is checked here because the emulator ignores it, so 
only a unit test can see
+   * it carried over.
+   */
+  @Test
+  public void testCreateBucketIsRoutedToV2WithProjectPrivateAcls() throws 
IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    // This is the bucket that GcpOptions.tryCreateDefaultBucketWithPrefix 
builds, plus a storage
+    // class.
+    Bucket bucket =
+        new Bucket()
+            .setName("bucket")
+            .setLocation("us-central1")
+            .setStorageClass("NEARLINE")
+            .setSoftDeletePolicy(new 
Bucket.SoftDeletePolicy().setRetentionDurationSeconds(0L));
+
+    gcsUtil.createBucket("a-project", bucket);
+
+    verify(mockDelegateV2)
+        .createBucket(
+            "a-project",
+            BucketInfo.newBuilder("bucket")
+                .setLocation("us-central1")
+                .setStorageClass(StorageClass.NEARLINE)
+                .setSoftDeletePolicy(
+                    BucketInfo.SoftDeletePolicy.newBuilder()
+                        .setRetentionDuration(java.time.Duration.ZERO)
+                        .build())
+                .build(),
+            BucketTargetOption.predefinedAcl(PredefinedAcl.PROJECT_PRIVATE),
+            
BucketTargetOption.predefinedDefaultObjectAcl(PredefinedAcl.PROJECT_PRIVATE));
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testRemoveBucketIsRoutedToV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+
+    gcsUtil.removeBucket(new Bucket().setName("bucket"));
+
+    verify(mockDelegateV2).removeBucket(BucketInfo.of("bucket"));
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testBucketOwnerIsRoutedToV2BucketProject() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+    when(mockDelegateV2.bucketProject(path)).thenReturn(123L);
+
+    assertEquals(123L, gcsUtil.bucketOwner(path));
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testGetObjectIsRoutedToV2AndKeepsAllFields() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+    com.google.cloud.storage.Blob blob = mockBlob("bucket", "object");
+    when(blob.getSize()).thenReturn(42L);
+    when(blob.getGeneration()).thenReturn(7L);
+    when(blob.getMetageneration()).thenReturn(3L);
+    when(blob.getContentType()).thenReturn("text/csv");
+    when(blob.getContentEncoding()).thenReturn("gzip");
+    when(blob.getMd5()).thenReturn("md5==");
+    when(blob.getCrc32c()).thenReturn("crc==");
+    when(blob.getEtag()).thenReturn("etag");
+    when(blob.getUpdateTimeOffsetDateTime())
+        
.thenReturn(java.time.Instant.ofEpochMilli(1234L).atOffset(java.time.ZoneOffset.UTC));
+    when(blob.getCreateTimeOffsetDateTime())
+        
.thenReturn(java.time.Instant.ofEpochMilli(1000L).atOffset(java.time.ZoneOffset.UTC));
+    when(mockDelegateV2.getBlob(path)).thenReturn(blob);
+
+    StorageObject object = gcsUtil.getObject(path);
+
+    assertEquals("bucket", object.getBucket());
+    assertEquals("object", object.getName());
+    assertEquals(BigInteger.valueOf(42L), object.getSize());
+    assertEquals(Long.valueOf(7L), object.getGeneration());
+    assertEquals(Long.valueOf(3L), object.getMetageneration());
+    assertEquals("text/csv", object.getContentType());
+    assertEquals("gzip", object.getContentEncoding());
+    assertEquals("md5==", object.getMd5Hash());
+    assertEquals("crc==", object.getCrc32c());
+    assertEquals("etag", object.getEtag());
+    assertEquals(1234L, object.getUpdated().getValue());
+    assertEquals(1000L, object.getTimeCreated().getValue());
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  /** Fields a blob may not carry (e.g. when fetched with a field mask) are 
left unset. */
+  @Test
+  public void testGetObjectLeavesMissingFieldsUnsetForV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+    com.google.cloud.storage.Blob blob = mockBlob("bucket", "object");
+    // Mockito would otherwise answer 0 for the boxed size.
+    when(blob.getSize()).thenReturn(null);
+    when(mockDelegateV2.getBlob(path)).thenReturn(blob);
+
+    StorageObject object = gcsUtil.getObject(path);
+
+    assertEquals("object", object.getName());
+    assertNull(object.getSize());
+    assertNull(object.getUpdated());
+    assertNull(object.getTimeCreated());
+  }
+
+  @Test
+  public void testCreateWithTypeIsRoutedToV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+
+    gcsUtil.create(path, "text/plain");
+    gcsUtil.create(path, "text/plain", 1024);
+
+    verify(mockDelegateV2)
+        .create(path, 
GcsUtilV1.CreateOptions.builder().setContentType("text/plain").build());
+    verify(mockDelegateV2)
+        .create(
+            path,
+            GcsUtilV1.CreateOptions.builder()
+                .setContentType("text/plain")
+                .setUploadBufferSizeBytes(1024)
+                .build());
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  /** GcpOptions reads the soft delete policy of the temp bucket through this 
method. */
+  @Test
+  public void testGetBucketIsRoutedToV2AndConvertsTheBucket() throws 
IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+    com.google.cloud.storage.Bucket bucket = 
Mockito.mock(com.google.cloud.storage.Bucket.class);
+    when(bucket.getName()).thenReturn("bucket");
+    when(bucket.getLocation()).thenReturn("US-CENTRAL1");
+    when(bucket.getProject()).thenReturn(BigInteger.valueOf(123L));
+    when(bucket.getStorageClass()).thenReturn(StorageClass.NEARLINE);
+    when(bucket.getSoftDeletePolicy())
+        .thenReturn(
+            BucketInfo.SoftDeletePolicy.newBuilder()
+                .setRetentionDuration(java.time.Duration.ofDays(7))
+                .build());
+    when(mockDelegateV2.getBucket(path)).thenReturn(bucket);
+
+    Bucket converted = gcsUtil.getBucket(path);
+
+    assertNotNull(converted);
+    assertEquals("bucket", converted.getName());
+    assertEquals("US-CENTRAL1", converted.getLocation());
+    assertEquals(BigInteger.valueOf(123L), converted.getProjectNumber());
+    assertEquals("NEARLINE", converted.getStorageClass());
+    assertEquals(
+        Long.valueOf(java.time.Duration.ofDays(7).getSeconds()),
+        converted.getSoftDeletePolicy().getRetentionDurationSeconds());
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testGetBucketWithoutSoftDeletePolicyForV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+    com.google.cloud.storage.Bucket bucket = 
Mockito.mock(com.google.cloud.storage.Bucket.class);
+    when(bucket.getName()).thenReturn("bucket");
+    when(mockDelegateV2.getBucket(path)).thenReturn(bucket);
+
+    Bucket converted = gcsUtil.getBucket(path);
+
+    assertNotNull(converted);
+    assertNull(converted.getSoftDeletePolicy());
+    assertNull(converted.getStorageClass());
+  }
+
+  /** Without the use_gcsutil_v2 experiment, the V2-only methods fail rather 
than fall back. */
+  @Test
+  public void testV2OnlyMethodsFailWithoutV2() {
+    GcsUtil gcsUtil = gcsOptionsWithTestCredential().getGcsUtil();
+    assertNull(gcsUtil.delegateV2);
+    GcsUtilV1 v1 = Mockito.mock(GcsUtilV1.class);
+    gcsUtil.delegate = v1;
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+    List<GcsPath> paths = ImmutableList.of(path);
+
+    List<ThrowingRunnable> calls =
+        ImmutableList.of(
+            () -> gcsUtil.getBlob(path),
+            () -> gcsUtil.getBlobs(paths),
+            () -> gcsUtil.listBlobs("bucket", "prefix", null),
+            () -> gcsUtil.listBlobs("bucket", "prefix", null, "/"),
+            () -> gcsUtil.openV2(path),
+            () -> gcsUtil.createV2(path, 
GcsUtil.CreateOptions.builder().build()),
+            () -> gcsUtil.createBucket(BucketInfo.of("bucket")),
+            () -> gcsUtil.getBucketWithOptions(path),
+            () -> gcsUtil.removeBucket(BucketInfo.of("bucket")),
+            () -> gcsUtil.copyV2(paths, paths),
+            () -> gcsUtil.copy(paths, paths, 
GcsUtilV2.OverwriteStrategy.ALWAYS_OVERWRITE),
+            () -> gcsUtil.renameV2(paths, paths),
+            () ->
+                gcsUtil.rename(
+                    paths,
+                    paths,
+                    GcsUtilV2.MissingStrategy.FAIL_IF_MISSING,
+                    GcsUtilV2.OverwriteStrategy.ALWAYS_OVERWRITE),
+            () -> gcsUtil.removeV2(paths),
+            () -> gcsUtil.remove(paths, 
GcsUtilV2.MissingStrategy.FAIL_IF_MISSING));
+
+    for (ThrowingRunnable call : calls) {
+      IOException e = assertThrows(IOException.class, call);
+      assertEquals("GcsUtil V2 not initialized.", e.getMessage());
+    }
+    Mockito.verifyNoInteractions(v1);
+  }
+
+  @Test
+  public void testExpandIsRoutedToV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath pattern = GcsPath.fromUri("gs://bucket/prefix/*");
+    List<GcsPath> expanded = 
ImmutableList.of(GcsPath.fromUri("gs://bucket/prefix/a"));
+    when(mockDelegateV2.expand(pattern)).thenReturn(expanded);
+
+    assertSame(expanded, gcsUtil.expand(pattern));
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testFileSizeIsRoutedToV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+    when(mockDelegateV2.fileSize(path)).thenReturn(42L);
+
+    assertEquals(42L, gcsUtil.fileSize(path));
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  /**
+   * Only the routing of the delimiter overload is checked. The page 
conversion is covered by {@link
+   * #testListObjectsIsRoutedToV2AndConvertsAPage}.
+   */
+  @Test
+  public void testListObjectsWithDelimiterIsRoutedToV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    com.google.cloud.storage.Blob object = mockBlob("bucket", "prefix/object");
+    @SuppressWarnings("unchecked")
+    Page<com.google.cloud.storage.Blob> page = Mockito.mock(Page.class);
+    when(page.getValues()).thenReturn(ImmutableList.of(object));
+    when(mockDelegateV2.listBlobs("bucket", "prefix/", "token", 
"/")).thenReturn(page);
+
+    Objects objects = gcsUtil.listObjects("bucket", "prefix/", "token", "/");
+
+    assertEquals("prefix/object", objects.getItems().get(0).getName());
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testOpenIsRoutedToV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+    SeekableByteChannel channel = Mockito.mock(SeekableByteChannel.class);
+    when(mockDelegateV2.open(path)).thenReturn(channel);
+
+    assertSame(channel, gcsUtil.open(path));
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testVerifyBucketAccessibleIsRoutedToV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath path = GcsPath.fromUri("gs://bucket/object");
+
+    gcsUtil.verifyBucketAccessible(path);
+
+    verify(mockDelegateV2).verifyBucketAccessible(path);
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  @Test
+  public void testBucketAccessibleIsRoutedToV2() throws IOException {
+    GcsUtil gcsUtil = gcsUtilRoutingToV2();
+    GcsPath accessible = GcsPath.fromUri("gs://accessible/object");
+    GcsPath inaccessible = GcsPath.fromUri("gs://inaccessible/object");
+    when(mockDelegateV2.bucketAccessible(accessible)).thenReturn(true);
+    when(mockDelegateV2.bucketAccessible(inaccessible)).thenReturn(false);
+
+    assertTrue(gcsUtil.bucketAccessible(accessible));
+    assertFalse(gcsUtil.bucketAccessible(inaccessible));
+    Mockito.verifyNoMoreInteractions(mockDelegate);
+  }
+
+  // The tests below exercise a real GcsUtilV2 delegate whose java-storage 
client is mocked, to
+  // cover behavior that GcsUtilV2 must share with GcsUtilV1.
+  // TODO: Move these to a parity test that runs against both delegates.
+
+  /**
+   * Returns a {@link GcsUtil} backed by a real {@link GcsUtilV2} that issues 
every call to {@code
+   * storage}. Performance metrics are off by default, so the per-operation 
clients of {@link
+   * GcsUtilV2#storageWithHttpMetrics} resolve to this one as well.
+   */
+  private GcsUtil gcsUtilWithV2Storage(com.google.cloud.storage.Storage 
storage) {
+    GcsOptions options = gcsOptionsWithTestCredential();
+    options.setProject("my_project");
+    GcsUtil gcsUtil = options.getGcsUtil();
+    GcsUtilV2 delegateV2 = Mockito.spy(new GcsUtilV2(options));
+    Mockito.doReturn(storage).when(delegateV2).storage();
+    gcsUtil.delegateV2 = delegateV2;
+    return gcsUtil;
+  }
+
+  /**
+   * Mirrors {@link #testGCSReadMetricsIsSet} for V2: opening a missing object 
records a {@code
+   * not_found} request, even though V2 detects it through a lookup rather 
than a {@link
+   * StorageException}.
+   */
+  @Test
+  public void testV2OpenMissingObjectRecordsNotFoundMetric() {
+    com.google.cloud.storage.Storage storage = 
Mockito.mock(com.google.cloud.storage.Storage.class);
+    // An unstubbed get() returns null, which is how java-storage reports a 
missing object.
+    GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+    assertThrows(
+        FileNotFoundException.class,
+        () -> gcsUtil.open(GcsPath.fromComponents("testbucket", 
"testobject")));
+
+    verifyMetricWasSet("my_project", "testbucket", "GcsGet", "not_found", 1);
+    verifyMetricWasSet("my_project", "testbucket", "GcsGet", "ok", 0);
+  }
+
+  /**
+   * An upload is finalized when its channel is closed, so a failed 
precondition surfaces there. V2
+   * must report it as an {@link IOException}, as V1 does, rather than an 
unchecked {@link
+   * StorageException}.
+   */
+  @Test
+  public void testV2WriteChannelCloseTranslatesStorageException() throws 
IOException {
+    com.google.cloud.storage.Storage storage = 
Mockito.mock(com.google.cloud.storage.Storage.class);
+    WriteChannel writer = Mockito.mock(WriteChannel.class);
+    StorageException preconditionFailed = new StorageException(412, 
"Precondition Failed");
+    Mockito.doThrow(preconditionFailed).when(writer).close();
+    when(storage.writer(any(com.google.cloud.storage.BlobInfo.class), 
any())).thenReturn(writer);
+    GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+    WritableByteChannel channel =
+        gcsUtil.create(
+            GcsPath.fromComponents("testbucket", "testobject"),
+            CreateOptions.builder().setExpectFileToNotExist(true).build());
+
+    IOException thrown = assertThrows(IOException.class, channel::close);
+    assertSame(preconditionFailed, thrown.getCause());
+  }
+
+  /** Mirrors {@link #testBucketDoesNotExist} for V2. */
+  @Test
+  public void testV2BucketAccessibleIsFalseWhenBucketDoesNotExist() throws 
IOException {
+    com.google.cloud.storage.Storage storage = 
Mockito.mock(com.google.cloud.storage.Storage.class);
+    // An unstubbed get() returns null, which is how java-storage reports a 
missing bucket.
+    GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+    assertFalse(gcsUtil.bucketAccessible(GcsPath.fromComponents("testbucket", 
"testobject")));
+  }
+
+  /** Mirrors {@link #testBucketDoesNotExistBecauseOfAccessError} for V2. */
+  @Test
+  public void testV2BucketAccessibleIsFalseWhenAccessIsDenied() throws 
IOException {
+    com.google.cloud.storage.Storage storage = 
Mockito.mock(com.google.cloud.storage.Storage.class);
+    when(storage.get(Mockito.eq("testbucket"), any(BucketGetOption.class)))
+        .thenThrow(new StorageException(403, "Forbidden"));
+    GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+    assertFalse(gcsUtil.bucketAccessible(GcsPath.fromComponents("testbucket", 
"testobject")));
+  }
+
+  /**
+   * Any other failure (e.g. a 5xx) says nothing about whether the bucket is 
accessible, so it must
+   * propagate rather than be reported as an inaccessible bucket, as V1 does.
+   */
+  @Test
+  public void testV2BucketAccessiblePropagatesOtherFailures() {
+    com.google.cloud.storage.Storage storage = 
Mockito.mock(com.google.cloud.storage.Storage.class);
+    StorageException serverError = new StorageException(503, "Service 
Unavailable");
+    when(storage.get(Mockito.eq("testbucket"), 
any(BucketGetOption.class))).thenThrow(serverError);
+    GcsUtil gcsUtil = gcsUtilWithV2Storage(storage);
+
+    IOException thrown =
+        assertThrows(
+            IOException.class,
+            () -> 
gcsUtil.bucketAccessible(GcsPath.fromComponents("testbucket", "testobject")));
+    assertSame(serverError, thrown.getCause());
+  }
 }
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
index 2b9e55b98b8..8bd79860cc2 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
@@ -30,6 +30,7 @@ import com.google.api.client.testing.http.MockHttpTransport;
 import com.google.api.client.testing.http.MockLowLevelHttpResponse;
 import com.google.auth.Credentials;
 import com.google.cloud.NoCredentials;
+import com.google.cloud.hadoop.util.AsyncWriteChannelOptions;
 import com.google.cloud.http.HttpTransportOptions;
 import com.google.cloud.storage.Storage;
 import com.google.cloud.storage.StorageOptions;
@@ -415,4 +416,15 @@ public class GcsUtilV2Test {
     assertNull(gcsCounter(other, "gcs_http_read_wire_bytes_received"));
     assertNull(gcsCounter(other, "gcs_http_write_wire_bytes_sent"));
   }
+
+  /**
+   * The chunk size decides how many requests a write costs, and the two 
clients do not default to
+   * the same one, so a drift here is a silent throughput regression rather 
than a test failure.
+   */
+  @Test
+  public void testDefaultUploadChunkSizeMatchesV1() {
+    assertEquals(
+        AsyncWriteChannelOptions.DEFAULT.getUploadChunkSize(),
+        GcsUtilV2.DEFAULT_UPLOAD_CHUNK_SIZE_BYTES);
+  }
 }

Reply via email to