Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
from google.cloud._storage_v2.services.storage.transports.base import (
DEFAULT_CLIENT_INFO,
)
from google.cloud.storage import __version__
from google.cloud.storage import __version__, _opentelemetry_metrics

_DEFAULT_HOST = "storage.googleapis.com"

Expand Down Expand Up @@ -63,6 +63,19 @@ class AsyncGrpcClient:
:param attempt_direct_path:
(Optional) Whether to attempt to use DirectPath for gRPC connections.
Defaults to ``True``.

:type enable_metrics: bool or None
:param enable_metrics:
(Optional, Experimental) Whether to enable OpenTelemetry metrics. If
None, falls back to the GCP_STORAGE_PYTHON_ENABLE_OTEL_METRICS
environment variable, or False if unset.

:type enable_debug_metrics: bool or None
:param enable_debug_metrics:
(Optional, Experimental) Whether to enable debug OpenTelemetry
metrics. If None, falls back to the
GCP_STORAGE_PYTHON_ENABLE_OTEL_DEBUG_METRICS environment variable, or False
if unset.
"""

def __init__(
Expand All @@ -72,7 +85,12 @@ def __init__(
client_options=None,
*,
attempt_direct_path=True,
enable_metrics=None,
enable_debug_metrics=None,
):
self._enable_metrics = enable_metrics
self._enable_debug_metrics = enable_debug_metrics

if isinstance(credentials, auth_credentials.AnonymousCredentials):
if client_options is None or client_options.api_endpoint is None:
raise ValueError(
Expand All @@ -99,6 +117,18 @@ def __init__(
attempt_direct_path=attempt_direct_path,
)

@property
def metrics_enabled(self) -> bool:
"""Returns True if metrics recording is active for this client."""
return _opentelemetry_metrics.is_metrics_enabled(self._enable_metrics)

@property
def debug_metrics_enabled(self) -> bool:
"""Returns True if debug metrics recording is active for this client."""
return _opentelemetry_metrics.is_debug_metrics_enabled(
self._enable_debug_metrics
)

def _create_anonymous_client(self, client_options, credentials):
channel = grpc.aio.insecure_channel(client_options.api_endpoint)
transport = storage_v2.services.storage.transports.StorageGrpcAsyncIOTransport(
Expand Down
30 changes: 30 additions & 0 deletions packages/google-cloud-storage/google/cloud/storage/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
from google.cloud.client import ClientWithProject
from google.cloud.exceptions import NotFound

from google.cloud.storage import _opentelemetry_metrics
from google.cloud.storage._bucket_metadata_cache import BucketMetadataCache
from google.cloud.storage._helpers import (
_DEFAULT_SCHEME,
Expand Down Expand Up @@ -126,6 +127,19 @@ class Client(ClientWithProject):
(Optional) An API key. Mutually exclusive with any other credentials.
This parameter is an alias for setting `client_options.api_key` and
will supercede any api key set in the `client_options` parameter.

:type enable_metrics: bool or None
:param enable_metrics:
(Optional, Experimental) Whether to enable OpenTelemetry metrics. If
None, falls back to the GCP_STORAGE_PYTHON_ENABLE_OTEL_METRICS
environment variable, or False if unset.

:type enable_debug_metrics: bool or None
:param enable_debug_metrics:
(Optional, Experimental) Whether to enable debug OpenTelemetry
metrics. If None, falls back to the
GCP_STORAGE_PYTHON_ENABLE_OTEL_DEBUG_METRICS environment variable, or False
if unset.
"""

SCOPE = (
Expand All @@ -146,6 +160,8 @@ def __init__(
extra_headers={},
*,
api_key=None,
enable_metrics=None,
enable_debug_metrics=None,
):
self._base_connection = None

Expand Down Expand Up @@ -293,6 +309,20 @@ def __init__(
self._connection = connection
self._batch_stack = _LocalStack()
self._bucket_metadata_cache = BucketMetadataCache(self)
self._enable_metrics = enable_metrics
self._enable_debug_metrics = enable_debug_metrics

@property
def metrics_enabled(self) -> bool:
"""Returns True if metrics recording is active for this client."""
return _opentelemetry_metrics.is_metrics_enabled(self._enable_metrics)

@property
def debug_metrics_enabled(self) -> bool:
"""Returns True if debug metrics recording is active for this client."""
return _opentelemetry_metrics.is_debug_metrics_enabled(
self._enable_debug_metrics
)

def close(self):
"""Close the client and clear any cached metadata or active connections."""
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
from google.cloud.client import ClientWithProject

from google.cloud import _storage_v2 as storage_v2
from google.cloud.storage import _opentelemetry_metrics

_marker = object()

Expand Down Expand Up @@ -58,6 +59,19 @@ class GrpcClient(ClientWithProject):
This provides a direct, unproxied connection to GCS for lower latency
and higher throughput, and is highly recommended when running on Google
Cloud infrastructure. Defaults to ``True``.

:type enable_metrics: bool or None
:param enable_metrics:
(Optional, Experimental) Whether to enable OpenTelemetry metrics. If
None, falls back to the GCP_STORAGE_PYTHON_ENABLE_OTEL_METRICS
environment variable, or False if unset.

:type enable_debug_metrics: bool or None
:param enable_debug_metrics:
(Optional, Experimental) Whether to enable debug OpenTelemetry
metrics. If None, falls back to the
GCP_STORAGE_PYTHON_ENABLE_OTEL_DEBUG_METRICS environment variable, or False
if unset.
"""

def __init__(
Expand All @@ -69,9 +83,14 @@ def __init__(
*,
api_key=None,
attempt_direct_path=True,
enable_metrics=None,
enable_debug_metrics=None,
):
super(GrpcClient, self).__init__(project=project, credentials=credentials)

self._enable_metrics = enable_metrics
self._enable_debug_metrics = enable_debug_metrics

if isinstance(client_options, dict):
if api_key:
client_options["api_key"] = api_key
Expand All @@ -87,6 +106,18 @@ def __init__(
attempt_direct_path=attempt_direct_path,
)

@property
def metrics_enabled(self) -> bool:
"""Returns True if metrics recording is active for this client."""
return _opentelemetry_metrics.is_metrics_enabled(self._enable_metrics)

@property
def debug_metrics_enabled(self) -> bool:
"""Returns True if debug metrics recording is active for this client."""
return _opentelemetry_metrics.is_debug_metrics_enabled(
self._enable_debug_metrics
)

def _create_gapic_client(
self,
credentials=None,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1383,6 +1383,8 @@ def _reduce_client(cl):
client_info = cl._initial_client_info
client_options = cl._initial_client_options
extra_headers = getattr(cl, "_extra_headers", {})
enable_metrics = getattr(cl, "_enable_metrics", None)
enable_debug_metrics = getattr(cl, "_enable_debug_metrics", None)

return _LazyClient, (
client_object_id,
Expand All @@ -1392,6 +1394,8 @@ def _reduce_client(cl):
client_info,
client_options,
extra_headers,
enable_metrics,
enable_debug_metrics,
)


Expand Down Expand Up @@ -1461,6 +1465,11 @@ def __new__(cls, id, *args, **kwargs):
if cached_client:
return cached_client
else:
if len(args) > 6:
kwargs.setdefault("enable_metrics", args[6])
if len(args) > 7:
kwargs.setdefault("enable_debug_metrics", args[7])
args = args[:6]
cached_client = Client(*args, **kwargs)
_cached_clients[id] = cached_client
return cached_client
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,45 @@ def test_constructor_disables_directpath(self, mock_async_storage_client):
mock_channel = mock_transport_cls.create_channel.return_value
mock_transport_cls.assert_called_once_with(channel=mock_channel)

@mock.patch("google.cloud._storage_v2.StorageAsyncClient")
def test_metrics_properties(self, mock_async_storage_client):
from google.cloud.storage import _opentelemetry_metrics

mock_transport_cls = mock.MagicMock()
mock_async_storage_client.get_transport_class.return_value = mock_transport_cls
mock_creds = _make_credentials()

client = async_grpc_client.AsyncGrpcClient(
credentials=mock_creds,
enable_metrics=True,
enable_debug_metrics=True,
)

with mock.patch.multiple(
_opentelemetry_metrics,
HAS_OPENTELEMETRY_METRICS=True,
_ENABLE_METRICS_DEV_GATE=True,
):
assert client.metrics_enabled is True
assert client.debug_metrics_enabled is True

client._enable_metrics = False
assert client.metrics_enabled is False
assert client.debug_metrics_enabled is True

client._enable_debug_metrics = False
assert client.debug_metrics_enabled is False

with mock.patch.multiple(
_opentelemetry_metrics,
HAS_OPENTELEMETRY_METRICS=True,
_ENABLE_METRICS_DEV_GATE=False,
):
client._enable_metrics = True
client._enable_debug_metrics = True
assert client.metrics_enabled is False
assert client.debug_metrics_enabled is False

@mock.patch("google.cloud._storage_v2.StorageAsyncClient")
def test_grpc_client_property(self, mock_grpc_gapic_client):
# Arrange
Expand Down
80 changes: 80 additions & 0 deletions packages/google-cloud-storage/tests/unit/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -259,6 +259,86 @@ def test_ctor_w_universe_domain_and_matched_credentials(self):
self.assertEqual(client.api_endpoint, expected_api_endpoint)
self.assertEqual(client.universe_domain, universe_domain)

def test_ctor_w_enable_metrics(self):
PROJECT = "PROJECT"
credentials = _make_credentials()
client = self._make_one(
project=PROJECT, credentials=credentials, enable_metrics=True
)
self.assertTrue(client._enable_metrics)

def test_ctor_w_enable_debug_metrics(self):
PROJECT = "PROJECT"
credentials = _make_credentials()
client = self._make_one(
project=PROJECT, credentials=credentials, enable_debug_metrics=True
)
self.assertTrue(client._enable_debug_metrics)

def test_client_metrics_enabled_property(self):
from google.cloud.storage import _opentelemetry_metrics

PROJECT = "PROJECT"
credentials = _make_credentials()
client = self._make_one(
project=PROJECT, credentials=credentials, enable_metrics=True
)

with mock.patch.multiple(
_opentelemetry_metrics,
HAS_OPENTELEMETRY_METRICS=True,
_ENABLE_METRICS_DEV_GATE=True,
):
self.assertTrue(client.metrics_enabled)

with mock.patch.multiple(
_opentelemetry_metrics,
HAS_OPENTELEMETRY_METRICS=True,
_ENABLE_METRICS_DEV_GATE=False,
):
self.assertFalse(client.metrics_enabled)

def test_client_debug_metrics_enabled_property(self):
from google.cloud.storage import _opentelemetry_metrics

PROJECT = "PROJECT"
credentials = _make_credentials()
client = self._make_one(
project=PROJECT,
credentials=credentials,
enable_metrics=True,
enable_debug_metrics=True,
)

with mock.patch.multiple(
_opentelemetry_metrics,
HAS_OPENTELEMETRY_METRICS=True,
_ENABLE_METRICS_DEV_GATE=True,
):
self.assertTrue(client.debug_metrics_enabled)

with mock.patch.multiple(
_opentelemetry_metrics,
HAS_OPENTELEMETRY_METRICS=True,
_ENABLE_METRICS_DEV_GATE=False,
):
self.assertFalse(client.debug_metrics_enabled)

# Verify that debug metrics are independent when base metrics are disabled
client_disabled_base = self._make_one(
project=PROJECT,
credentials=credentials,
enable_metrics=False,
enable_debug_metrics=True,
)
with mock.patch.multiple(
_opentelemetry_metrics,
HAS_OPENTELEMETRY_METRICS=True,
_ENABLE_METRICS_DEV_GATE=True,
):
self.assertFalse(client_disabled_base.metrics_enabled)
self.assertTrue(client_disabled_base.debug_metrics_enabled)

def test_ctor_w_universe_domain_and_mismatched_credentials(self):
PROJECT = "PROJECT"
universe_domain = "example.com"
Expand Down
40 changes: 40 additions & 0 deletions packages/google-cloud-storage/tests/unit/test_grpc_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -197,3 +197,43 @@ def test_constructor_with_api_key_and_dict_options(
client_info=None,
client_options=expected_options,
)

@mock.patch("google.cloud.storage.grpc_client.ClientWithProject")
@mock.patch("google.cloud._storage_v2.StorageClient")
def test_metrics_properties(self, mock_storage_client, mock_base_client):
from google.cloud.storage import _opentelemetry_metrics

mock_creds = _make_credentials()
mock_base_client.return_value._credentials = mock_creds

client = grpc_client.GrpcClient(
project="test-project",
credentials=mock_creds,
enable_metrics=True,
enable_debug_metrics=True,
)

with mock.patch.multiple(
_opentelemetry_metrics,
HAS_OPENTELEMETRY_METRICS=True,
_ENABLE_METRICS_DEV_GATE=True,
):
self.assertTrue(client.metrics_enabled)
self.assertTrue(client.debug_metrics_enabled)

client._enable_metrics = False
self.assertFalse(client.metrics_enabled)
self.assertTrue(client.debug_metrics_enabled)

client._enable_debug_metrics = False
self.assertFalse(client.debug_metrics_enabled)

with mock.patch.multiple(
_opentelemetry_metrics,
HAS_OPENTELEMETRY_METRICS=True,
_ENABLE_METRICS_DEV_GATE=False,
):
client._enable_metrics = True
client._enable_debug_metrics = True
self.assertFalse(client.metrics_enabled)
self.assertFalse(client.debug_metrics_enabled)
Loading
Loading