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 720326829a7 Fill gaps between GcsUtilV2 and GcsUtilV1 (#40244)
720326829a7 is described below

commit 720326829a71b16b3553a246305fbf9031a70d3d
Author: Shunping Huang <[email protected]>
AuthorDate: Thu Sep 24 22:03:02 2026 -0400

    Fill gaps between GcsUtilV2 and GcsUtilV1 (#40244)
    
    * Fix a leaking channel in gcsutilv2.
    
    * Configure GcsUtilV2's storage client from pipeline options
    
    GcsUtilV2 left credentials to the application default chain, silently
    ignoring --gcpCredentialFactoryClass, impersonation and explicit service
    account keys that GcsUtilV1 honors. A custom --gcsEndpoint was ignored
    outright.
    
    Pass the pipeline's credentials through, mapping a null credential (an
    explicit opt-out) to NoCredentials, and apply gcsEndpoint as the client
    host.
    
    * Emit GCS metrics from GcsUtilV2 read and write paths
    
    GcsUtilV2 returned raw channels, so enabling use_gcsutil_v2 silently
    dropped every GCS metric that GcsUtilV1 reports: the per-bucket byte
    counters from GcsCountersOptions, the gcs_http_*_wire_bytes counters
    gated on gcsPerformanceMetrics, and the api_request_count
    ServiceCallMetric.
    
    Wrap the channels returned by open() and create() in the existing
    Counting{Seekable,Writable}ByteChannel helpers and report a
    ServiceCallMetric for GcsGet and GcsInsert, matching V1's labels.
    
    * Implement GCS Http request metrics for GcsUtilV2
    
    * Add tests for GCS HTTP request metrics for GcsUtilV2
    
    * Log which GcsUtil delegate is in use.
    
    * Fix spotbug finding
    
    * Fix another spotbug finding
    
    * Address reviewer comments to add metric and configuration tests.
---
 .../org/apache/beam/gradle/BeamModulePlugin.groovy |   1 +
 .../google-cloud-platform-core/build.gradle        |   1 +
 .../beam/sdk/extensions/gcp/util/GcsUtil.java      |  10 +-
 .../beam/sdk/extensions/gcp/util/GcsUtilV2.java    | 268 ++++++++++++-
 .../gcp/util/GcsUtilParameterizedIT.java           | 155 ++++++++
 .../sdk/extensions/gcp/util/GcsUtilV2Test.java     | 418 +++++++++++++++++++++
 6 files changed, 840 insertions(+), 13 deletions(-)

diff --git 
a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy 
b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
index 88c998ed910..0c34e1529e8 100644
--- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
+++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
@@ -768,6 +768,7 @@ class BeamModulePlugin implements Plugin<Project> {
         google_cloud_bigtable_emulator              : 
"com.google.cloud:google-cloud-bigtable-emulator", // 
google_cloud_platform_libraries_bom sets version
         google_cloud_core                           : 
"com.google.cloud:google-cloud-core", // google_cloud_platform_libraries_bom 
sets version
         google_cloud_core_grpc                      : 
"com.google.cloud:google-cloud-core-grpc", // 
google_cloud_platform_libraries_bom sets version
+        google_cloud_core_http                      : 
"com.google.cloud:google-cloud-core-http", // 
google_cloud_platform_libraries_bom sets version
         google_cloud_datacatalog_v1beta1            : 
"com.google.cloud:google-cloud-datacatalog", // 
google_cloud_platform_libraries_bom sets version
         google_cloud_dataflow_java_proto_library_all: 
"com.google.cloud.dataflow:google-cloud-dataflow-java-proto-library-all:0.5.160304",
         google_cloud_datastore_v1_proto_client      : 
"com.google.cloud.datastore:datastore-v1-proto-client:3.4.0",   // 
[bomupgrader] sets version
diff --git a/sdks/java/extensions/google-cloud-platform-core/build.gradle 
b/sdks/java/extensions/google-cloud-platform-core/build.gradle
index 75e0c508727..9f2b475a713 100644
--- a/sdks/java/extensions/google-cloud-platform-core/build.gradle
+++ b/sdks/java/extensions/google-cloud-platform-core/build.gradle
@@ -44,6 +44,7 @@ dependencies {
   implementation library.java.google_auth_library_oauth2_http
   implementation library.java.google_api_client
   implementation library.java.google_cloud_core
+  implementation library.java.google_cloud_core_http
   implementation library.java.google_cloud_storage
   implementation library.java.bigdataoss_gcsio
   implementation library.java.bigdataoss_util
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 070cf74d7c1..9e2d82ae1eb 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
@@ -47,8 +47,12 @@ import org.apache.beam.sdk.options.PipelineOptions;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
 import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Sets;
 import org.checkerframework.checker.nullness.qual.Nullable;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 public class GcsUtil {
+  private static final Logger LOG = LoggerFactory.getLogger(GcsUtil.class);
+
   /**
    * Namespace for every GCS metric. The namespace is dropped when Dataflow 
exports counters to
    * Cloud Monitoring, so the layer is carried by the metric name instead: 
{@code gcs_http_*} for
@@ -119,8 +123,12 @@ public class GcsUtil {
     this.delegate = new GcsUtilV1.GcsUtilFactory().create(options);
     if (ExperimentalOptions.hasExperiment(options, "use_gcsutil_v2")) {
       this.delegateV2 = new GcsUtilV2.GcsUtilFactory().create(options);
+      // INFO only for V2, which is opt-in. V1 is still the default for every 
pipeline,
+      // so logging it at INFO would be noise.
+      LOG.info("Using GcsUtilV2 (java-storage) for GCS operations.");
     } else {
       this.delegateV2 = null;
+      LOG.debug("Using GcsUtilV1 (gcsio) for GCS operations.");
     }
   }
 
@@ -297,7 +305,7 @@ public class GcsUtil {
 
   public WritableByteChannel create(GcsPath path, CreateOptions options) 
throws IOException {
     if (delegateV2 != null) {
-      delegateV2.create(path, options.delegate);
+      return delegateV2.create(path, options.delegate);
     }
     return delegate.create(path, options.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 bbdab0982b5..9d223f81050 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
@@ -21,10 +21,16 @@ import static 
org.apache.beam.sdk.io.FileSystemUtils.wildcardToRegexp;
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
 import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
 
+import com.google.api.client.http.HttpRequestInitializer;
 import com.google.api.gax.paging.Page;
+import com.google.auth.Credentials;
 import com.google.auto.value.AutoValue;
+import com.google.cloud.NoCredentials;
 import com.google.cloud.ReadChannel;
+import com.google.cloud.ServiceOptions;
+import com.google.cloud.TransportOptions;
 import com.google.cloud.WriteChannel;
+import com.google.cloud.http.HttpTransportOptions;
 import com.google.cloud.storage.Blob;
 import com.google.cloud.storage.BlobId;
 import com.google.cloud.storage.BlobInfo;
@@ -47,6 +53,8 @@ import com.google.cloud.storage.StorageException;
 import com.google.cloud.storage.StorageOptions;
 import java.io.FileNotFoundException;
 import java.io.IOException;
+import java.net.MalformedURLException;
+import java.net.URL;
 import java.nio.ByteBuffer;
 import java.nio.channels.SeekableByteChannel;
 import java.nio.channels.WritableByteChannel;
@@ -54,11 +62,25 @@ import java.nio.file.AccessDeniedException;
 import java.nio.file.FileAlreadyExistsException;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.function.Consumer;
 import java.util.regex.Pattern;
+import org.apache.beam.runners.core.metrics.GcpResourceIdentifiers;
+import org.apache.beam.runners.core.metrics.MonitoringInfoConstants;
+import org.apache.beam.runners.core.metrics.ServiceCallMetric;
 import org.apache.beam.sdk.extensions.gcp.options.GcpOptions;
 import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
+import 
org.apache.beam.sdk.extensions.gcp.util.channels.CountingSeekableByteChannel;
+import 
org.apache.beam.sdk.extensions.gcp.util.channels.CountingWritableByteChannel;
 import org.apache.beam.sdk.extensions.gcp.util.gcsfs.GcsPath;
+import org.apache.beam.sdk.metrics.Counter;
+import org.apache.beam.sdk.metrics.MetricName;
+import org.apache.beam.sdk.metrics.Metrics;
+import org.apache.beam.sdk.metrics.MetricsContainer;
+import org.apache.beam.sdk.metrics.MetricsEnvironment;
 import org.apache.beam.sdk.options.DefaultValueFactory;
 import org.apache.beam.sdk.options.PipelineOptions;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
@@ -82,6 +104,12 @@ class GcsUtilV2 {
 
   private final @Nullable Integer uploadBufferSizeBytes;
 
+  private final GcsUtilV1.GcsCountersOptions gcsCountersOptions;
+
+  private final boolean gcsPerformanceMetrics;
+
+  private final @Nullable String projectId;
+
   /** Maximum number of items to retrieve per Objects.List request. */
   private static final long MAX_LIST_BLOBS_PER_CALL = 1024;
 
@@ -95,9 +123,204 @@ class GcsUtilV2 {
   private static final long MEGABYTES_COPIED_PER_CHUNK = 2048L;
 
   GcsUtilV2(PipelineOptions options) {
-    String projectId = options.as(GcpOptions.class).getProject();
-    storage = 
StorageOptions.newBuilder().setProjectId(projectId).build().getService();
-    uploadBufferSizeBytes = 
options.as(GcsOptions.class).getGcsUploadBufferSizeBytes();
+    GcsOptions gcsOptions = options.as(GcsOptions.class);
+    this.projectId = options.as(GcpOptions.class).getProject();
+    StorageOptions.Builder storageOptionsBuilder =
+        StorageOptions.newBuilder().setProjectId(this.projectId);
+
+    // Use the pipeline's configured credentials rather than falling back to 
application default
+    // credentials, so that --gcpCredentialFactoryClass, impersonation and 
explicit service account
+    // keys are honored. A null credential means the pipeline opted out of 
authentication
+    // (e.g. NoopCredentialFactory), which maps to NoCredentials for this 
client.
+    Credentials credentials = gcsOptions.getGcpCredential();
+    storageOptionsBuilder.setCredentials(
+        credentials != null ? credentials : NoCredentials.getInstance());
+
+    // GcsOptions#getGcsEndpoint may carry a service path (as the JSON client 
in Transport expects),
+    // but this client derives its own path, so only the root is applicable 
here.
+    String endpoint = gcsOptions.getGcsEndpoint();
+    if (endpoint != null) {
+      storageOptionsBuilder.setHost(rootUrlOf(endpoint));
+    }
+
+    storage = storageOptionsBuilder.build().getService();
+    uploadBufferSizeBytes = gcsOptions.getGcsUploadBufferSizeBytes();
+    this.gcsCountersOptions =
+        GcsUtilV1.GcsCountersOptions.create(
+            gcsOptions.getEnableBucketReadMetricCounter()
+                ? gcsOptions.getGcsReadCounterPrefix()
+                : null,
+            gcsOptions.getEnableBucketWriteMetricCounter()
+                ? gcsOptions.getGcsWriteCounterPrefix()
+                : null);
+    this.gcsPerformanceMetrics = 
Boolean.TRUE.equals(gcsOptions.getGcsPerformanceMetrics());
+  }
+
+  /**
+   * Creates an integer consumer that updates the counter identified by a 
prefix and a bucket name.
+   */
+  private static Consumer<Integer> createCounterConsumer(String 
counterNamePrefix, String bucket) {
+    return Metrics.counter(GcsUtil.class, String.format("%s_%s", 
counterNamePrefix, bucket))::inc;
+  }
+
+  /** Returns the {@link MetricsContainer} to attribute wire-byte counters to, 
if enabled. */
+  @VisibleForTesting
+  @Nullable MetricsContainer performanceMetricsContainer() {
+    return gcsPerformanceMetrics ? MetricsEnvironment.getCurrentContainer() : 
null;
+  }
+
+  @VisibleForTesting
+  WritableByteChannel wrapInCounting(
+      WritableByteChannel writableByteChannel,
+      String bucket,
+      @Nullable MetricsContainer container) {
+    Consumer<Integer> writeConsumer =
+        Optional.ofNullable(gcsCountersOptions.getWriteCounterPrefix())
+            .map(prefix -> createCounterConsumer(prefix, bucket))
+            .orElse(null);
+
+    if (this.gcsPerformanceMetrics && container != null) {
+      Counter perfWriteCounter =
+          container.getCounter(
+              MetricName.named(GcsUtil.METRIC_NAMESPACE, 
"gcs_http_write_wire_bytes_sent"));
+      Consumer<Integer> perfConsumer = perfWriteCounter::inc;
+      writeConsumer = writeConsumer == null ? perfConsumer : 
writeConsumer.andThen(perfConsumer);
+    }
+
+    if (writeConsumer == null) {
+      return writableByteChannel;
+    }
+    return new CountingWritableByteChannel(writableByteChannel, writeConsumer);
+  }
+
+  @VisibleForTesting
+  SeekableByteChannel wrapInCounting(
+      SeekableByteChannel seekableByteChannel,
+      String bucket,
+      @Nullable MetricsContainer container) {
+    Consumer<Integer> readConsumer =
+        Optional.ofNullable(gcsCountersOptions.getReadCounterPrefix())
+            .map(prefix -> createCounterConsumer(prefix, bucket))
+            .orElse(null);
+
+    if (this.gcsPerformanceMetrics && container != null) {
+      Counter perfReadCounter =
+          container.getCounter(
+              MetricName.named(GcsUtil.METRIC_NAMESPACE, 
"gcs_http_read_wire_bytes_received"));
+      Consumer<Integer> perfConsumer = perfReadCounter::inc;
+      readConsumer = readConsumer == null ? perfConsumer : 
readConsumer.andThen(perfConsumer);
+    }
+
+    if (readConsumer == null) {
+      return seekableByteChannel;
+    }
+    return CountingSeekableByteChannel.createWithBytesReadConsumer(
+        seekableByteChannel, readConsumer);
+  }
+
+  /** Builds the API request metric for {@code method} (e.g. {@code GcsGet}) 
on {@code bucket}. */
+  private ServiceCallMetric serviceCallMetric(String method, String bucket) {
+    HashMap<String, String> baseLabels = new HashMap<>();
+    baseLabels.put(MonitoringInfoConstants.Labels.PTRANSFORM, "");
+    baseLabels.put(MonitoringInfoConstants.Labels.SERVICE, "Storage");
+    baseLabels.put(MonitoringInfoConstants.Labels.METHOD, method);
+    baseLabels.put(
+        MonitoringInfoConstants.Labels.RESOURCE, 
GcpResourceIdentifiers.cloudStorageBucket(bucket));
+    baseLabels.put(MonitoringInfoConstants.Labels.GCS_PROJECT_ID, 
String.valueOf(projectId));
+    baseLabels.put(MonitoringInfoConstants.Labels.GCS_BUCKET, bucket);
+    return new 
ServiceCallMetric(MonitoringInfoConstants.Urns.API_REQUEST_COUNT, baseLabels);
+  }
+
+  /**
+   * {@link HttpTransportOptions} that wraps the request initializer so that 
HTTP-level counters
+   * (request counts, request shape, and status classes) are incremented 
against a pre-bound {@link
+   * MetricsContainer}.
+   *
+   * <p>The container is bound eagerly rather than resolved per request 
because requests may execute
+   * on background threads, where {@link 
MetricsEnvironment#getCurrentContainer} would not resolve
+   * to the step that initiated the operation.
+   */
+  private static class MetricsHttpTransportOptions extends 
HttpTransportOptions {
+    private static final long serialVersionUID = 1L;
+
+    // Not serializable, and only meaningful in the process that created it.
+    private final transient @Nullable MetricsContainer container;
+    private final boolean isWrite;
+
+    MetricsHttpTransportOptions(
+        HttpTransportOptions base, @Nullable MetricsContainer container, 
boolean isWrite) {
+      super(base.toBuilder());
+      this.container = container;
+      this.isWrite = isWrite;
+    }
+
+    @Override
+    public HttpRequestInitializer getHttpRequestInitializer(ServiceOptions<?, 
?> serviceOptions) {
+      // withMetricsContainer returns the delegate unchanged when the 
container is null, which is
+      // also the case after deserialization.
+      return Transport.withMetricsContainer(
+          super.getHttpRequestInitializer(serviceOptions), container, isWrite);
+    }
+
+    @Override
+    public boolean equals(@Nullable Object obj) {
+      if (this == obj) {
+        return true;
+      }
+      if (!(obj instanceof MetricsHttpTransportOptions)) {
+        return false;
+      }
+      if (!super.equals(obj)) {
+        return false;
+      }
+      MetricsHttpTransportOptions other = (MetricsHttpTransportOptions) obj;
+      return isWrite == other.isWrite && Objects.equals(container, 
other.container);
+    }
+
+    @Override
+    public int hashCode() {
+      return Objects.hash(super.hashCode(), container, isWrite);
+    }
+  }
+
+  /**
+   * Returns a client whose HTTP requests are counted against {@code 
container}, or the shared
+   * client when HTTP metrics are not being collected.
+   *
+   * <p>A distinct client is required because the interceptors are installed 
on the transport, which
+   * is fixed when the client is built. This mirrors {@code GcsUtilV1}, which 
likewise builds a
+   * scoped client per operation while performance metrics are enabled.
+   */
+  @VisibleForTesting
+  Storage storageWithHttpMetrics(@Nullable MetricsContainer container, boolean 
isWrite) {
+    if (container == null) {
+      return storage;
+    }
+    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 options.toBuilder()
+        .setTransportOptions(
+            new MetricsHttpTransportOptions(
+                (HttpTransportOptions) transportOptions, container, isWrite))
+        .build()
+        .getService();
+  }
+
+  /** Returns the {@code scheme://host[:port]} prefix of {@code endpoint}, 
discarding any path. */
+  private static String rootUrlOf(String endpoint) {
+    try {
+      URL url = new URL(endpoint);
+      return url.getProtocol()
+          + "://"
+          + url.getHost()
+          + (url.getPort() > 0 ? ":" + url.getPort() : "");
+    } catch (MalformedURLException e) {
+      throw new IllegalArgumentException("Invalid gcsEndpoint URL: " + 
endpoint, e);
+    }
   }
 
   // AccessDeniedException/FileAlreadyExistsException permit a null "other" 
argument, and these
@@ -123,8 +346,14 @@ class GcsUtilV2 {
   }
 
   public Blob getBlob(GcsPath gcsPath, BlobGetOption... options) throws 
IOException {
+    return getBlob(storage, gcsPath, options);
+  }
+
+  /** As {@link #getBlob(GcsPath, BlobGetOption...)}, but issued through a 
specific client. */
+  private Blob getBlob(Storage client, GcsPath gcsPath, BlobGetOption... 
options)
+      throws IOException {
     try {
-      Blob blob = storage.get(gcsPath.getBucket(), gcsPath.getObject(), 
options);
+      Blob blob = client.get(gcsPath.getBucket(), gcsPath.getObject(), 
options);
       if (blob == null) {
         throw new FileNotFoundException(
             String.format("The specified file does not exist: %s", 
gcsPath.toString()));
@@ -561,11 +790,21 @@ class GcsUtilV2 {
 
   public SeekableByteChannel open(GcsPath path, BlobSourceOption... 
sourceOptions)
       throws IOException {
-    Blob blob = getBlob(path, BlobGetOption.fields(BlobField.SIZE));
-    ReadChannel reader = blob.getStorage().reader(blob.getBlobId(), 
sourceOptions);
-    // disable internal buffering, and make the channel non-blocking
-    reader.setChunkSize(0);
-    return new GcsSeekableByteChannel(reader, blob.getSize());
+    ServiceCallMetric serviceCallMetric = serviceCallMetric("GcsGet", 
path.getBucket());
+    MetricsContainer container = performanceMetricsContainer();
+    try {
+      Storage client = storageWithHttpMetrics(container, false);
+      Blob blob = getBlob(client, path, BlobGetOption.fields(BlobField.SIZE));
+      ReadChannel reader = client.reader(blob.getBlobId(), sourceOptions);
+      // disable internal buffering, and make the channel non-blocking
+      reader.setChunkSize(0);
+      serviceCallMetric.call("ok");
+      return wrapInCounting(
+          new GcsSeekableByteChannel(reader, blob.getSize()), 
path.getBucket(), container);
+    } catch (StorageException e) {
+      serviceCallMetric.call(e.getCode());
+      throw translateStorageException(path, e);
+    }
   }
 
   /** A bridge that allows a GCS WriteChannel to behave as a 
WritableByteChannel. */
@@ -601,6 +840,8 @@ class GcsUtilV2 {
   public WritableByteChannel create(
       GcsPath path, GcsUtilV1.CreateOptions options, BlobWriteOption... 
writeOptions)
       throws IOException {
+    ServiceCallMetric serviceCallMetric = serviceCallMetric("GcsInsert", 
path.getBucket());
+    MetricsContainer container = performanceMetricsContainer();
     try {
       // Define the metadata for the new object
       BlobInfo.Builder builder = BlobInfo.newBuilder(path.getBucket(), 
path.getObject());
@@ -611,13 +852,14 @@ class GcsUtilV2 {
 
       BlobInfo blobInfo = builder.build();
 
+      Storage client = storageWithHttpMetrics(container, true);
       List<BlobWriteOption> writeOptionList = new 
ArrayList<>(Arrays.asList(writeOptions));
       if (options.getExpectFileToNotExist()) {
         writeOptionList.add(BlobWriteOption.doesNotExist());
       } else {
         // We do not merge this check with the getExpectFileToNotExist() 
branch above
         // because we don't want to always make the storage.get() RPC call.
-        Blob blob = storage.get(path.getBucket(), path.getObject());
+        Blob blob = client.get(path.getBucket(), path.getObject());
         if (blob == null) {
           writeOptionList.add(BlobWriteOption.doesNotExist());
         } else {
@@ -626,7 +868,7 @@ class GcsUtilV2 {
       }
       // Open a WriteChannel from the storage service
       WriteChannel writer =
-          storage.writer(blobInfo, writeOptionList.toArray(new 
BlobWriteOption[0]));
+          client.writer(blobInfo, writeOptionList.toArray(new 
BlobWriteOption[0]));
       Integer uploadBufferSizeBytes =
           options.getUploadBufferSizeBytes() != null
               ? options.getUploadBufferSizeBytes()
@@ -635,10 +877,12 @@ class GcsUtilV2 {
         writer.setChunkSize(uploadBufferSizeBytes);
       }
 
+      serviceCallMetric.call("ok");
       // Return the bridge wrapper
-      return new GcsWritableByteChannel(writer, path);
+      return wrapInCounting(new GcsWritableByteChannel(writer, path), 
path.getBucket(), container);
 
     } catch (StorageException e) {
+      serviceCallMetric.call(e.getCode());
       throw translateStorageException(path, e);
     }
   }
diff --git 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilParameterizedIT.java
 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilParameterizedIT.java
index 5759bb10a65..17752337ce9 100644
--- 
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilParameterizedIT.java
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilParameterizedIT.java
@@ -20,6 +20,8 @@ package org.apache.beam.sdk.extensions.gcp.util;
 import static org.junit.Assert.assertArrayEquals;
 import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
 import static org.junit.Assert.assertThrows;
 import static org.junit.Assert.assertTrue;
 
@@ -45,13 +47,22 @@ import java.security.NoSuchAlgorithmException;
 import java.util.Arrays;
 import java.util.Collections;
 import java.util.List;
+import java.util.Map;
 import java.util.stream.Collectors;
+import org.apache.beam.runners.core.metrics.CounterCell;
+import org.apache.beam.runners.core.metrics.GcpResourceIdentifiers;
+import org.apache.beam.runners.core.metrics.MetricUpdates.MetricUpdate;
+import org.apache.beam.runners.core.metrics.MetricsContainerImpl;
+import org.apache.beam.runners.core.metrics.MonitoringInfoConstants;
+import org.apache.beam.runners.core.metrics.MonitoringInfoMetricName;
 import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
 import org.apache.beam.sdk.extensions.gcp.util.GcsUtil.CreateOptions;
 import org.apache.beam.sdk.extensions.gcp.util.GcsUtilV2.MissingStrategy;
 import org.apache.beam.sdk.extensions.gcp.util.GcsUtilV2.OverwriteStrategy;
 import org.apache.beam.sdk.extensions.gcp.util.gcsfs.GcsPath;
 import org.apache.beam.sdk.io.fs.MoveOptions;
+import org.apache.beam.sdk.metrics.MetricName;
+import org.apache.beam.sdk.metrics.MetricsEnvironment;
 import org.apache.beam.sdk.options.ExperimentalOptions;
 import org.apache.beam.sdk.testing.TestPipeline;
 import org.apache.beam.sdk.testing.TestPipelineOptions;
@@ -75,6 +86,9 @@ import org.junit.runners.Parameterized.Parameters;
 @Category(UsesKms.class)
 public class GcsUtilParameterizedIT {
 
+  private static final String READ_COUNTER_PREFIX = "it_read_bytes";
+  private static final String WRITE_COUNTER_PREFIX = "it_write_bytes";
+
   @Parameters(name = "{0}")
   public static Iterable<String> data() {
     return Arrays.asList("use_gcsutil_v1", "use_gcsutil_v2");
@@ -683,4 +697,145 @@ public class GcsUtilParameterizedIT {
       tearDownTestBucketHelper(bucketName);
     }
   }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Metrics parity: the same assertions must hold for both GcsUtilV1 and 
GcsUtilV2.
+  // 
---------------------------------------------------------------------------------------------
+
+  /** Returns a {@link GcsUtil} with every GCS metric flag turned on. */
+  private GcsUtil gcsUtilWithAllMetrics() {
+    GcsOptions gcsOptions = options.as(GcsOptions.class);
+    gcsOptions.setGcsPerformanceMetrics(true);
+    gcsOptions.setEnableBucketReadMetricCounter(true);
+    gcsOptions.setEnableBucketWriteMetricCounter(true);
+    gcsOptions.setGcsReadCounterPrefix(READ_COUNTER_PREFIX);
+    gcsOptions.setGcsWriteCounterPrefix(WRITE_COUNTER_PREFIX);
+    // Built directly, as getGcsUtil() returns the instance cached in setUp().
+    return new GcsUtil(gcsOptions);
+  }
+
+  private static long counter(MetricsContainerImpl container, MetricName name) 
{
+    CounterCell cell = container.tryGetCounter(name);
+    assertNotNull("counter " + name + " was not reported", cell);
+    return cell.getCumulative();
+  }
+
+  private static long gcsCounter(MetricsContainerImpl container, String name) {
+    return counter(container, MetricName.named(GcsUtil.METRIC_NAMESPACE, 
name));
+  }
+
+  private static long bucketCounter(MetricsContainerImpl container, String 
prefix, String bucket) {
+    return counter(container, MetricName.named(GcsUtil.class, prefix + "_" + 
bucket));
+  }
+
+  /**
+   * Sums the API request counter for {@code method} and {@code status} on 
{@code bucket}.
+   *
+   * <p>The {@code GCS_PROJECT_ID} label is deliberately not matched: V1 
reports the project of its
+   * gcsio options (which is never set, so it reports {@code "null"}), while 
V2 reports the
+   * pipeline's project.
+   */
+  private static long apiRequestCount(
+      MetricsContainerImpl container, String method, String status, String 
bucket) {
+    long total = 0;
+    for (MetricUpdate<Long> update : 
container.getCumulative().counterUpdates()) {
+      MetricName name = update.getKey().metricName();
+      if (!(name instanceof MonitoringInfoMetricName)) {
+        continue;
+      }
+      MonitoringInfoMetricName miName = (MonitoringInfoMetricName) name;
+      Map<String, String> labels = miName.getLabels();
+      if 
(MonitoringInfoConstants.Urns.API_REQUEST_COUNT.equals(miName.getUrn())
+          && 
"Storage".equals(labels.get(MonitoringInfoConstants.Labels.SERVICE))
+          && method.equals(labels.get(MonitoringInfoConstants.Labels.METHOD))
+          && status.equals(labels.get(MonitoringInfoConstants.Labels.STATUS))
+          && 
bucket.equals(labels.get(MonitoringInfoConstants.Labels.GCS_BUCKET))
+          && GcpResourceIdentifiers.cloudStorageBucket(bucket)
+              .equals(labels.get(MonitoringInfoConstants.Labels.RESOURCE))) {
+        total += update.getUpdate();
+      }
+    }
+    return total;
+  }
+
+  @Test
+  public void testReadMetrics() throws IOException {
+    final String bucket = "apache-beam-samples";
+    final GcsPath gcsPath = GcsPath.fromComponents(bucket, 
"shakespeare/kinglear.txt");
+    final long expectedSize = 157283L;
+    GcsUtil metricsGcsUtil = gcsUtilWithAllMetrics();
+
+    MetricsContainerImpl container = new MetricsContainerImpl("step");
+    MetricsContainerImpl processWide = new MetricsContainerImpl(null);
+    MetricsEnvironment.setCurrentContainer(container);
+    MetricsEnvironment.setProcessWideContainer(processWide);
+    try {
+      try (SeekableByteChannel channel = metricsGcsUtil.open(gcsPath)) {
+        ByteBuffer buffer = ByteBuffer.allocate((int) expectedSize + 1024);
+        assertEquals(expectedSize, 
StorageChannelUtils.blockingFillFrom(buffer, channel));
+      }
+    } finally {
+      MetricsEnvironment.setCurrentContainer(null);
+      MetricsEnvironment.setProcessWideContainer(null);
+    }
+
+    // --enableBucketReadMetricCounter
+    assertEquals(expectedSize, bucketCounter(container, READ_COUNTER_PREFIX, 
bucket));
+    // --gcsPerformanceMetrics: wire bytes and HTTP counters
+    assertEquals(expectedSize, gcsCounter(container, 
"gcs_http_read_wire_bytes_received"));
+    assertTrue(gcsCounter(container, "gcs_http_read_request_count") >= 1);
+    assertTrue(gcsCounter(container, "gcs_http_read_status_2xx") >= 1);
+    assertEquals(
+        gcsCounter(container, "gcs_http_read_request_count"),
+        gcsCounter(container, "gcs_http_read_request_count_ranged")
+            + gcsCounter(container, "gcs_http_read_request_count_unbounded")
+            + gcsCounter(container, "gcs_http_read_request_count_other"));
+    // A read must not produce write-side counters.
+    assertNull(
+        container.tryGetCounter(
+            MetricName.named(GcsUtil.METRIC_NAMESPACE, 
"gcs_http_write_request_count")));
+    // API request metric
+    assertEquals(1, apiRequestCount(processWide, "GcsGet", "ok", bucket));
+  }
+
+  @Test
+  public void testWriteMetrics() throws IOException {
+    final String bucket =
+        "apache-beam-temp-metrics-" + 
java.util.UUID.randomUUID().toString().substring(0, 8);
+    final GcsPath targetPath = GcsPath.fromComponents(bucket, 
"test-object.txt");
+    final byte[] content = "Hello, GCS 
metrics!".getBytes(StandardCharsets.UTF_8);
+    GcsUtil metricsGcsUtil = gcsUtilWithAllMetrics();
+
+    MetricsContainerImpl container = new MetricsContainerImpl("step");
+    MetricsContainerImpl processWide = new MetricsContainerImpl(null);
+    try {
+      createTestBucketHelper(bucket, false);
+
+      MetricsEnvironment.setCurrentContainer(container);
+      MetricsEnvironment.setProcessWideContainer(processWide);
+      try (WritableByteChannel writer =
+          metricsGcsUtil.create(
+              targetPath, 
CreateOptions.builder().setExpectFileToNotExist(true).build())) {
+        writer.write(ByteBuffer.wrap(content));
+      } finally {
+        MetricsEnvironment.setCurrentContainer(null);
+        MetricsEnvironment.setProcessWideContainer(null);
+      }
+    } finally {
+      tearDownTestBucketHelper(bucket);
+    }
+
+    // --enableBucketWriteMetricCounter
+    assertEquals(content.length, bucketCounter(container, 
WRITE_COUNTER_PREFIX, bucket));
+    // --gcsPerformanceMetrics: wire bytes and HTTP counters
+    assertEquals(content.length, gcsCounter(container, 
"gcs_http_write_wire_bytes_sent"));
+    assertTrue(gcsCounter(container, "gcs_http_write_request_count") >= 1);
+    assertTrue(gcsCounter(container, "gcs_http_write_status_2xx") >= 1);
+    // A write must not produce read-side counters.
+    assertNull(
+        container.tryGetCounter(
+            MetricName.named(GcsUtil.METRIC_NAMESPACE, 
"gcs_http_read_request_count")));
+    // API request metric
+    assertEquals(1, apiRequestCount(processWide, "GcsInsert", "ok", bucket));
+  }
 }
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
new file mode 100644
index 00000000000..2b9e55b98b8
--- /dev/null
+++ 
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV2Test.java
@@ -0,0 +1,418 @@
+/*
+ * 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.beam.sdk.extensions.gcp.util;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotSame;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertThrows;
+
+import com.google.api.client.http.GenericUrl;
+import com.google.api.client.http.HttpRequestInitializer;
+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.http.HttpTransportOptions;
+import com.google.cloud.storage.Storage;
+import com.google.cloud.storage.StorageOptions;
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.nio.channels.Channels;
+import java.nio.channels.SeekableByteChannel;
+import java.nio.channels.WritableByteChannel;
+import 
org.apache.beam.repackaged.core.org.apache.commons.compress.utils.SeekableInMemoryByteChannel;
+import org.apache.beam.runners.core.metrics.CounterCell;
+import org.apache.beam.runners.core.metrics.MetricsContainerImpl;
+import org.apache.beam.sdk.extensions.gcp.auth.NoopCredentialFactory;
+import org.apache.beam.sdk.extensions.gcp.auth.TestCredential;
+import org.apache.beam.sdk.extensions.gcp.options.GcsOptions;
+import org.apache.beam.sdk.metrics.MetricName;
+import org.apache.beam.sdk.metrics.MetricsEnvironment;
+import org.apache.beam.sdk.options.PipelineOptionsFactory;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.junit.After;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/**
+ * Unit tests for {@link GcsUtilV2}.
+ *
+ * <p>Covers two things without making any request to GCS:
+ *
+ * <ul>
+ *   <li>The storage client is configured from the pipeline options (project, 
credentials,
+ *       endpoint), matching what {@link GcsUtilV1} honors.
+ *   <li>HTTP metrics and byte counters are installed only when their flags 
ask for them, and only
+ *       when there is a container to report into. The counting logic itself 
is covered by {@link
+ *       TransportTest}.
+ * </ul>
+ *
+ * <p>End-to-end parity with {@link GcsUtilV1} against real GCS is covered by 
{@link
+ * GcsUtilParameterizedIT}.
+ */
+@RunWith(JUnit4.class)
+public class GcsUtilV2Test {
+
+  private static final String BUCKET = "test-bucket";
+  private static final String READ_PREFIX = "test_read_prefix";
+  private static final String WRITE_PREFIX = "test_write_prefix";
+  private static final byte[] PAYLOAD = "some_bytes".getBytes(UTF_8);
+
+  @After
+  public void tearDown() {
+    MetricsEnvironment.setCurrentContainer(null);
+  }
+
+  private static GcsOptions gcsOptions() {
+    GcsOptions options = PipelineOptionsFactory.as(GcsOptions.class);
+    // Avoid resolving application default credentials; no request leaves the 
process.
+    options.setGcpCredential(new TestCredential());
+    options.setProject("test-project");
+    return options;
+  }
+
+  private static GcsUtilV2 gcsUtilV2(boolean performanceMetrics) {
+    GcsOptions options = gcsOptions();
+    options.setGcsPerformanceMetrics(performanceMetrics);
+    return new GcsUtilV2(options);
+  }
+
+  private static GcsUtilV2 gcsUtilV2(boolean performanceMetrics, boolean 
bucketCounters) {
+    GcsOptions options = gcsOptions();
+    options.setGcsPerformanceMetrics(performanceMetrics);
+    options.setEnableBucketReadMetricCounter(bucketCounters);
+    options.setEnableBucketWriteMetricCounter(bucketCounters);
+    options.setGcsReadCounterPrefix(READ_PREFIX);
+    options.setGcsWriteCounterPrefix(WRITE_PREFIX);
+    return new GcsUtilV2(options);
+  }
+
+  /** The shared (unscoped) storage client of {@code gcsUtil}. */
+  private static StorageOptions storageOptionsOf(GcsUtilV2 gcsUtil) {
+    return gcsUtil.storageWithHttpMetrics(null, false).getOptions();
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Configuration
+  // 
---------------------------------------------------------------------------------------------
+
+  @Test
+  public void testStorageClientUsesPipelineProject() {
+    assertEquals("test-project", 
storageOptionsOf(gcsUtilV2(false)).getProjectId());
+  }
+
+  @Test
+  public void testStorageClientUsesPipelineCredentials() {
+    GcsOptions options = gcsOptions();
+    Credentials credentials = options.getGcpCredential();
+
+    assertSame(credentials, storageOptionsOf(new 
GcsUtilV2(options)).getCredentials());
+  }
+
+  @Test
+  public void testNullCredentialMapsToNoCredentials() {
+    GcsOptions options = gcsOptions();
+    options.setCredentialFactoryClass(NoopCredentialFactory.class);
+    options.setGcpCredential(null);
+
+    assertSame(
+        NoCredentials.getInstance(), storageOptionsOf(new 
GcsUtilV2(options)).getCredentials());
+  }
+
+  @Test
+  public void testDefaultHostWithoutGcsEndpoint() {
+    String defaultHost =
+        StorageOptions.newBuilder()
+            .setProjectId("test-project")
+            .setCredentials(NoCredentials.getInstance())
+            .build()
+            .getHost();
+
+    assertEquals(defaultHost, storageOptionsOf(gcsUtilV2(false)).getHost());
+  }
+
+  /** Mirrors {@code GcsUtilTest#testGcsEndpoint}: only the root of the 
endpoint applies to V2. */
+  @Test
+  public void testGcsEndpointRootIsUsedAsHost() {
+    GcsOptions options = gcsOptions();
+    options.setGcsEndpoint("http://localhost:4443/storage/v1/";);
+
+    assertEquals("http://localhost:4443";, storageOptionsOf(new 
GcsUtilV2(options)).getHost());
+  }
+
+  @Test
+  public void testGcsEndpointWithoutPort() {
+    GcsOptions options = gcsOptions();
+    options.setGcsEndpoint("https://storage.example.com/storage/v1/";);
+
+    assertEquals("https://storage.example.com";, storageOptionsOf(new 
GcsUtilV2(options)).getHost());
+  }
+
+  @Test
+  public void testInvalidGcsEndpointIsRejected() {
+    GcsOptions options = gcsOptions();
+    options.setGcsEndpoint("not a url");
+
+    assertThrows(IllegalArgumentException.class, () -> new GcsUtilV2(options));
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // HTTP metrics binding
+  // 
---------------------------------------------------------------------------------------------
+
+  @Test
+  public void testSharedClientIsReusedWithoutAContainer() {
+    GcsUtilV2 gcsUtil = gcsUtilV2(true);
+
+    Storage read = gcsUtil.storageWithHttpMetrics(null, false);
+    Storage write = gcsUtil.storageWithHttpMetrics(null, true);
+
+    // With nothing to count, no per-operation client should be built.
+    assertSame(read, write);
+  }
+
+  @Test
+  public void testScopedClientIsBuiltForAContainer() {
+    GcsUtilV2 gcsUtil = gcsUtilV2(true);
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+
+    Storage shared = gcsUtil.storageWithHttpMetrics(null, false);
+    Storage scoped = gcsUtil.storageWithHttpMetrics(container, false);
+
+    assertNotSame(shared, scoped);
+    assertNotSame(
+        shared.getOptions().getTransportOptions(), 
scoped.getOptions().getTransportOptions());
+    // Everything other than the transport must carry over from the shared 
client.
+    assertEquals(shared.getOptions().getProjectId(), 
scoped.getOptions().getProjectId());
+    assertEquals(shared.getOptions().getHost(), scoped.getOptions().getHost());
+    assertSame(shared.getOptions().getCredentials(), 
scoped.getOptions().getCredentials());
+  }
+
+  private static long counter(MetricsContainerImpl container, String name) {
+    return container.getCounter(MetricName.named(GcsUtil.METRIC_NAMESPACE, 
name)).getCumulative();
+  }
+
+  /** Executes one request through {@code initializer} against a mock 
transport. */
+  private static void executeOneRequest(HttpRequestInitializer initializer) 
throws IOException {
+    MockHttpTransport transport =
+        new MockHttpTransport.Builder()
+            .setLowLevelHttpResponse(new 
MockLowLevelHttpResponse().setStatusCode(200))
+            .build();
+    transport
+        .createRequestFactory(initializer)
+        .buildGetRequest(new GenericUrl("https://storage.googleapis.com/test";))
+        .execute();
+  }
+
+  private static HttpRequestInitializer initializerOf(Storage client) {
+    return ((HttpTransportOptions) client.getOptions().getTransportOptions())
+        .getHttpRequestInitializer(client.getOptions());
+  }
+
+  /**
+   * The scoped client must hand out a request initializer that counts, since 
that is the only way
+   * the HTTP counters reach the container.
+   */
+  @Test
+  public void testScopedClientInitializerCountsRequests() throws IOException {
+    GcsUtilV2 gcsUtil = gcsUtilV2(true);
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+
+    executeOneRequest(initializerOf(gcsUtil.storageWithHttpMetrics(container, 
false)));
+
+    assertEquals(1, counter(container, "gcs_http_read_request_count"));
+    assertEquals(1, counter(container, "gcs_http_read_status_2xx"));
+  }
+
+  @Test
+  public void testWriteDirectionIsCountedSeparately() throws IOException {
+    GcsUtilV2 gcsUtil = gcsUtilV2(true);
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+
+    executeOneRequest(initializerOf(gcsUtil.storageWithHttpMetrics(container, 
true)));
+
+    assertEquals(1, counter(container, "gcs_http_write_request_count"));
+    assertEquals(0, counter(container, "gcs_http_read_request_count"));
+  }
+
+  /**
+   * The flag is the actual gate: with performance metrics disabled no 
container is resolved, so
+   * open() and create() pass null and no scoped client or HTTP counter is 
ever produced.
+   */
+  @Test
+  public void testPerformanceMetricsFlagGatesTheContainer() {
+    MetricsContainerImpl current = new MetricsContainerImpl(null);
+    MetricsEnvironment.setCurrentContainer(current);
+    try {
+      assertSame(current, gcsUtilV2(true).performanceMetricsContainer());
+      assertNull(gcsUtilV2(false).performanceMetricsContainer());
+    } finally {
+      MetricsEnvironment.setCurrentContainer(null);
+    }
+  }
+
+  // 
---------------------------------------------------------------------------------------------
+  // Byte counters (mirrors GcsUtilTest#testReadMetrics / #testWriteMetrics 
for V1)
+  // 
---------------------------------------------------------------------------------------------
+
+  private static @Nullable Long counterOrNull(MetricsContainerImpl container, 
MetricName name) {
+    CounterCell cell = container.tryGetCounter(name);
+    return cell == null ? null : cell.getCumulative();
+  }
+
+  private static @Nullable Long bucketCounter(MetricsContainerImpl container, 
String prefix) {
+    return counterOrNull(container, MetricName.named(GcsUtil.class, prefix + 
"_" + BUCKET));
+  }
+
+  private static @Nullable Long gcsCounter(MetricsContainerImpl container, 
String name) {
+    return counterOrNull(container, MetricName.named(GcsUtil.METRIC_NAMESPACE, 
name));
+  }
+
+  /** Reads {@link #PAYLOAD} through V2's read wrapper, bound to {@code 
bound}. */
+  private static void readPayload(GcsUtilV2 gcsUtil, @Nullable 
MetricsContainerImpl bound)
+      throws IOException {
+    try (SeekableByteChannel channel =
+        gcsUtil.wrapInCounting(new SeekableInMemoryByteChannel(PAYLOAD), 
BUCKET, bound)) {
+      assertEquals(PAYLOAD.length, 
channel.read(ByteBuffer.allocate(PAYLOAD.length)));
+    }
+  }
+
+  /** Writes {@link #PAYLOAD} through V2's write wrapper, bound to {@code 
bound}. */
+  private static void writePayload(GcsUtilV2 gcsUtil, @Nullable 
MetricsContainerImpl bound)
+      throws IOException {
+    try (WritableByteChannel channel =
+        gcsUtil.wrapInCounting(Channels.newChannel(new 
ByteArrayOutputStream()), BUCKET, bound)) {
+      assertEquals(PAYLOAD.length, channel.write(ByteBuffer.wrap(PAYLOAD)));
+    }
+  }
+
+  @Test
+  public void testReadCountersWhenAllEnabled() throws IOException {
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+    MetricsEnvironment.setCurrentContainer(container);
+
+    readPayload(gcsUtilV2(true, true), container);
+
+    assertEquals(Long.valueOf(PAYLOAD.length), bucketCounter(container, 
READ_PREFIX));
+    assertEquals(
+        Long.valueOf(PAYLOAD.length), gcsCounter(container, 
"gcs_http_read_wire_bytes_received"));
+    assertNull(bucketCounter(container, WRITE_PREFIX));
+    assertNull(gcsCounter(container, "gcs_http_write_wire_bytes_sent"));
+  }
+
+  @Test
+  public void testWriteCountersWhenAllEnabled() throws IOException {
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+    MetricsEnvironment.setCurrentContainer(container);
+
+    writePayload(gcsUtilV2(true, true), container);
+
+    assertEquals(Long.valueOf(PAYLOAD.length), bucketCounter(container, 
WRITE_PREFIX));
+    assertEquals(
+        Long.valueOf(PAYLOAD.length), gcsCounter(container, 
"gcs_http_write_wire_bytes_sent"));
+    assertNull(bucketCounter(container, READ_PREFIX));
+    assertNull(gcsCounter(container, "gcs_http_read_wire_bytes_received"));
+  }
+
+  @Test
+  public void testChannelsAreNotWrappedWhenAllDisabled() throws IOException {
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+    MetricsEnvironment.setCurrentContainer(container);
+    GcsUtilV2 gcsUtil = gcsUtilV2(false, false);
+    SeekableByteChannel readChannel = new SeekableInMemoryByteChannel(PAYLOAD);
+    WritableByteChannel writeChannel = Channels.newChannel(new 
ByteArrayOutputStream());
+
+    assertSame(readChannel, gcsUtil.wrapInCounting(readChannel, BUCKET, 
container));
+    assertSame(writeChannel, gcsUtil.wrapInCounting(writeChannel, BUCKET, 
container));
+  }
+
+  @Test
+  public void testOnlyBucketCountersWhenPerformanceMetricsDisabled() throws 
IOException {
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+    MetricsEnvironment.setCurrentContainer(container);
+    GcsUtilV2 gcsUtil = gcsUtilV2(false, true);
+
+    readPayload(gcsUtil, container);
+    writePayload(gcsUtil, container);
+
+    assertEquals(Long.valueOf(PAYLOAD.length), bucketCounter(container, 
READ_PREFIX));
+    assertEquals(Long.valueOf(PAYLOAD.length), bucketCounter(container, 
WRITE_PREFIX));
+    assertNull(gcsCounter(container, "gcs_http_read_wire_bytes_received"));
+    assertNull(gcsCounter(container, "gcs_http_write_wire_bytes_sent"));
+  }
+
+  @Test
+  public void testOnlyWireBytesWhenBucketCountersDisabled() throws IOException 
{
+    MetricsContainerImpl container = new MetricsContainerImpl(null);
+    MetricsEnvironment.setCurrentContainer(container);
+    GcsUtilV2 gcsUtil = gcsUtilV2(true, false);
+
+    readPayload(gcsUtil, container);
+    writePayload(gcsUtil, container);
+
+    assertNull(bucketCounter(container, READ_PREFIX));
+    assertNull(bucketCounter(container, WRITE_PREFIX));
+    assertEquals(
+        Long.valueOf(PAYLOAD.length), gcsCounter(container, 
"gcs_http_read_wire_bytes_received"));
+    assertEquals(
+        Long.valueOf(PAYLOAD.length), gcsCounter(container, 
"gcs_http_write_wire_bytes_sent"));
+  }
+
+  /** Without a bound container there is nowhere to report wire bytes, so 
nothing is wrapped. */
+  @Test
+  public void testWireBytesRequireABoundContainer() {
+    GcsUtilV2 gcsUtil = gcsUtilV2(true, false);
+    SeekableByteChannel readChannel = new SeekableInMemoryByteChannel(PAYLOAD);
+    WritableByteChannel writeChannel = Channels.newChannel(new 
ByteArrayOutputStream());
+
+    assertSame(readChannel, gcsUtil.wrapInCounting(readChannel, BUCKET, null));
+    assertSame(writeChannel, gcsUtil.wrapInCounting(writeChannel, BUCKET, 
null));
+  }
+
+  /**
+   * Wire bytes go to the container bound when the channel was created, even 
if a different
+   * container is current while the bytes are read or written.
+   */
+  @Test
+  public void testWireBytesAreAttributedToTheBoundContainer() throws 
IOException {
+    MetricsContainerImpl bound = new MetricsContainerImpl("bound");
+    MetricsContainerImpl other = new MetricsContainerImpl("other");
+    GcsUtilV2 gcsUtil = gcsUtilV2(true, false);
+    SeekableByteChannel readChannel =
+        gcsUtil.wrapInCounting(new SeekableInMemoryByteChannel(PAYLOAD), 
BUCKET, bound);
+    WritableByteChannel writeChannel =
+        gcsUtil.wrapInCounting(Channels.newChannel(new 
ByteArrayOutputStream()), BUCKET, bound);
+
+    MetricsEnvironment.setCurrentContainer(other);
+    readChannel.read(ByteBuffer.allocate(PAYLOAD.length));
+    writeChannel.write(ByteBuffer.wrap(PAYLOAD));
+    readChannel.close();
+    writeChannel.close();
+
+    assertEquals(
+        Long.valueOf(PAYLOAD.length), gcsCounter(bound, 
"gcs_http_read_wire_bytes_received"));
+    assertEquals(Long.valueOf(PAYLOAD.length), gcsCounter(bound, 
"gcs_http_write_wire_bytes_sent"));
+    assertNull(gcsCounter(other, "gcs_http_read_wire_bytes_received"));
+    assertNull(gcsCounter(other, "gcs_http_write_wire_bytes_sent"));
+  }
+}

Reply via email to