From 3d33f794267699a7d93a0825808691c00e447e73 Mon Sep 17 00:00:00 2001 From: Nidhi Nandwani Date: Sun, 20 Sep 2026 13:06:34 +0000 Subject: [PATCH 1/4] feat(storage): support DirectPath over Interconnect in GCS gRPC Add support for DirectPath over Cloud Interconnect (DP over GCI) across google.api_core.grpc_helpers, google.api_core.grpc_helpers_async, google.cloud.storage.grpc_client.GrpcClient, and google.cloud.storage.asyncio.async_grpc_client.AsyncGrpcClient. - Add attempt_direct_path_xds_over_interconnect option and GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT env override. - Synthesize standard TLS composite credentials instead of GCE ALTS when DirectPath over Interconnect is enabled. - Rewrite storage.googleapis.com to storage-direct.googleapis.com and append ?force-xds to google-c2p:/// target URIs. [Generated-by: AI] --- .../google/api_core/grpc_helpers.py | 73 ++++++++-- .../google/api_core/grpc_helpers_async.py | 30 +++- .../tests/asyncio/test_grpc_helpers_async.py | 92 ++++++++++++ .../tests/unit/test_grpc_helpers.py | 137 ++++++++++++++++++ .../storage/asyncio/async_grpc_client.py | 53 +++++-- .../google/cloud/storage/grpc_client.py | 65 ++++++++- .../unit/asyncio/test_async_grpc_client.py | 59 ++++++++ .../tests/unit/test_grpc_client.py | 97 +++++++++++++ 8 files changed, 579 insertions(+), 27 deletions(-) diff --git a/packages/google-api-core/google/api_core/grpc_helpers.py b/packages/google-api-core/google/api_core/grpc_helpers.py index 71ed528f2795..27fca814425f 100644 --- a/packages/google-api-core/google/api_core/grpc_helpers.py +++ b/packages/google-api-core/google/api_core/grpc_helpers.py @@ -16,6 +16,7 @@ import collections import functools +import os import warnings from typing import ( Callable, @@ -38,6 +39,21 @@ from google.api_core import exceptions, general_helpers +_DIRECT_PATH_INTERCONNECT_ENV = "GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT" + + +def _resolve_direct_path_interconnect( + attempt_direct_path_xds_over_interconnect: Optional[bool], +) -> bool: + """Resolves whether DirectPath over Interconnect is enabled.""" + env_val = os.environ.get(_DIRECT_PATH_INTERCONNECT_ENV) + if env_val == "true": + return True + if env_val == "false": + return False + return bool(attempt_direct_path_xds_over_interconnect) + + # The list of gRPC Callable interfaces that return iterators. _STREAM_WRAP_CLASSES = (grpc.UnaryStreamMultiCallable, grpc.StreamStreamMultiCallable) @@ -324,6 +340,7 @@ def create_channel( default_host=None, compression=None, attempt_direct_path: Optional[bool] = False, + attempt_direct_path_xds_over_interconnect: Optional[bool] = False, **kwargs, ): """Create a secure channel with credentials. @@ -374,6 +391,9 @@ def create_channel( `False` as the Service may not support Direct Path. - Using `ssl_credentials` with `attempt_direct_path` set to `True` will result in `ValueError` as this combination is not yet supported. + attempt_direct_path_xds_over_interconnect (Optional[bool]): If set, + DirectPath over Cloud Interconnect will be attempted using standard + TLS credentials and ``?force-xds`` C2P target resolution. kwargs: Additional key-word args passed to :func:`grpc.secure_channel`. @@ -382,16 +402,25 @@ def create_channel( grpc.Channel: The created channel. Raises: - google.api_core.DuplicateCredentialArgs: If both a credentials object and credentials_file are passed. - ValueError: If `ssl_credentials` is set and `attempt_direct_path` is set to `True`. + google.api_core.DuplicateCredentialArgs: If both a credentials object + and credentials_file are passed. + ValueError: If `ssl_credentials` is set and `attempt_direct_path` is + set to `True` without `attempt_direct_path_xds_over_interconnect`. """ + use_dp_interconnect = _resolve_direct_path_interconnect( + attempt_direct_path_xds_over_interconnect + ) + # If `ssl_credentials` is set and `attempt_direct_path` is set to `True`, - # raise ValueError as this is not yet supported. + # raise ValueError as this is not yet supported for GCE ALTS DirectPath. # See https://github.com/googleapis/python-api-core/issues/590 - if ssl_credentials and attempt_direct_path: + if ssl_credentials and attempt_direct_path and not use_dp_interconnect: raise ValueError("Using ssl_credentials with Direct Path is not supported") + if use_dp_interconnect and ssl_credentials is None: + ssl_credentials = grpc.ssl_channel_credentials() + composite_credentials = _create_composite_credentials( credentials=credentials, credentials_file=credentials_file, @@ -402,21 +431,31 @@ def create_channel( default_host=default_host, ) - if attempt_direct_path: - target = _modify_target_for_direct_path(target) + if attempt_direct_path or use_dp_interconnect: + target = _modify_target_for_direct_path( + target, + attempt_direct_path_xds_over_interconnect=use_dp_interconnect, + ) + elif "-direct." in target and not target.startswith("google-c2p:///"): + target = target.replace("-direct.", ".") return grpc.secure_channel( target, composite_credentials, compression=compression, **kwargs ) -def _modify_target_for_direct_path(target: str) -> str: +def _modify_target_for_direct_path( + target: str, + attempt_direct_path_xds_over_interconnect: Optional[bool] = False, +) -> str: """ Given a target, return a modified version which is compatible with Direct Path. Args: target (str): The target service address in the format 'hostname[:port]' or 'dns://hostname[:port]'. + attempt_direct_path_xds_over_interconnect (Optional[bool]): Whether to + append ``?force-xds`` for DirectPath over Cloud Interconnect. Returns: target (str): The target service address which is converted into a format compatible with Direct Path. @@ -434,9 +473,23 @@ def _modify_target_for_direct_path(target: str) -> str: direct_path_separator = ":///" if direct_path_separator not in target: - target_without_port = target.split(":")[0] - # Modify the target to use Direct Path by adding the `google-c2p:///` prefix - target = f"google-c2p{direct_path_separator}{target_without_port}" + if "?" in target: + host_part, query_part = target.split("?", 1) + target_without_port = host_part.split(":")[0] + target = ( + f"google-c2p{direct_path_separator}{target_without_port}?{query_part}" + ) + else: + target_without_port = target.split(":")[0] + # Modify the target to use Direct Path by adding the `google-c2p:///` prefix + target = f"google-c2p{direct_path_separator}{target_without_port}" + + if attempt_direct_path_xds_over_interconnect and target.startswith( + "google-c2p:///" + ): + if "force-xds" not in target: + separator = "&" if "?" in target else "?" + target = f"{target}{separator}force-xds" return target diff --git a/packages/google-api-core/google/api_core/grpc_helpers_async.py b/packages/google-api-core/google/api_core/grpc_helpers_async.py index d1f897901e7a..13940885b4c7 100644 --- a/packages/google-api-core/google/api_core/grpc_helpers_async.py +++ b/packages/google-api-core/google/api_core/grpc_helpers_async.py @@ -219,6 +219,7 @@ def create_channel( default_host=None, compression=None, attempt_direct_path: Optional[bool] = False, + attempt_direct_path_xds_over_interconnect: Optional[bool] = False, **kwargs, ): """Create an AsyncIO secure channel with credentials. @@ -270,6 +271,9 @@ def create_channel( `False` as the Service may not support Direct Path. - Using `ssl_credentials` with `attempt_direct_path` set to `True` will result in `ValueError` as this combination is not yet supported. + attempt_direct_path_xds_over_interconnect (Optional[bool]): If set, + DirectPath over Cloud Interconnect will be attempted using standard + TLS credentials and ``?force-xds`` C2P target resolution. kwargs: Additional key-word args passed to :func:`aio.secure_channel`. @@ -277,19 +281,28 @@ def create_channel( aio.Channel: The created channel. Raises: - google.api_core.DuplicateCredentialArgs: If both a credentials object and credentials_file are passed. - ValueError: If `ssl_credentials` is set and `attempt_direct_path` is set to `True`. + google.api_core.DuplicateCredentialArgs: If both a credentials object + and credentials_file are passed. + ValueError: If `ssl_credentials` is set and `attempt_direct_path` is + set to `True` without `attempt_direct_path_xds_over_interconnect`. """ if credentials_file is not None: warnings.warn(general_helpers._CREDENTIALS_FILE_WARNING, DeprecationWarning) + use_dp_interconnect = grpc_helpers._resolve_direct_path_interconnect( + attempt_direct_path_xds_over_interconnect + ) + # If `ssl_credentials` is set and `attempt_direct_path` is set to `True`, - # raise ValueError as this is not yet supported. + # raise ValueError as this is not yet supported for GCE ALTS DirectPath. # See https://github.com/googleapis/python-api-core/issues/590 - if ssl_credentials and attempt_direct_path: + if ssl_credentials and attempt_direct_path and not use_dp_interconnect: raise ValueError("Using ssl_credentials with Direct Path is not supported") + if use_dp_interconnect and ssl_credentials is None: + ssl_credentials = grpc.ssl_channel_credentials() + composite_credentials = grpc_helpers._create_composite_credentials( credentials=credentials, credentials_file=credentials_file, @@ -300,8 +313,13 @@ def create_channel( default_host=default_host, ) - if attempt_direct_path: - target = grpc_helpers._modify_target_for_direct_path(target) + if attempt_direct_path or use_dp_interconnect: + target = grpc_helpers._modify_target_for_direct_path( + target, + attempt_direct_path_xds_over_interconnect=use_dp_interconnect, + ) + elif "-direct." in target and not target.startswith("google-c2p:///"): + target = target.replace("-direct.", ".") return aio.secure_channel( target, composite_credentials, compression=compression, **kwargs diff --git a/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py b/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py index dcb09f18fea2..bab77fe765c4 100644 --- a/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py +++ b/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py @@ -743,3 +743,95 @@ async def test_fake_stream_unary_call(): await fake_call.wait_for_connection() response = await fake_call assert fake_call.response == response + + +@mock.patch("grpc.ssl_channel_credentials") +@mock.patch("grpc.composite_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.aio.secure_channel") +def test_create_channel_attempt_direct_path_xds_over_interconnect_default_ssl( + grpc_secure_channel, + google_auth_default, + composite_creds_call, + ssl_creds_call, +): + composite_creds = composite_creds_call.return_value + default_ssl_creds = ssl_creds_call.return_value + + channel = grpc_helpers_async.create_channel( + "storage-direct.googleapis.com:443", + attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=True, + ) + + assert channel is grpc_secure_channel.return_value + ssl_creds_call.assert_called_once_with() + assert composite_creds_call.call_args.args[0] is default_ssl_creds + grpc_secure_channel.assert_called_once_with( + "google-c2p:///storage-direct.googleapis.com?force-xds", + composite_creds, + compression=None, + ) + + +@mock.patch("grpc.composite_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.aio.secure_channel") +def test_create_channel_attempt_direct_path_xds_over_interconnect_custom_ssl( + grpc_secure_channel, + google_auth_default, + composite_creds_call, +): + composite_creds = composite_creds_call.return_value + custom_ssl_creds = mock.sentinel.custom_ssl_creds + + channel = grpc_helpers_async.create_channel( + "storage-direct.googleapis.com:443", + ssl_credentials=custom_ssl_creds, + attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=True, + ) + + assert channel is grpc_secure_channel.return_value + assert composite_creds_call.call_args.args[0] is custom_ssl_creds + grpc_secure_channel.assert_called_once_with( + "google-c2p:///storage-direct.googleapis.com?force-xds", + composite_creds, + compression=None, + ) + + +@mock.patch("grpc.compute_engine_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.aio.secure_channel") +def test_create_channel_interconnect_fallback( + grpc_secure_channel, + google_auth_default, + composite_creds_call, +): + composite_creds = composite_creds_call.return_value + + channel = grpc_helpers_async.create_channel( + "storage-direct.googleapis.com:443", + attempt_direct_path=False, + attempt_direct_path_xds_over_interconnect=False, + ) + + assert channel is grpc_secure_channel.return_value + grpc_secure_channel.assert_called_once_with( + "storage.googleapis.com:443", + composite_creds, + compression=None, + ) diff --git a/packages/google-api-core/tests/unit/test_grpc_helpers.py b/packages/google-api-core/tests/unit/test_grpc_helpers.py index 1a1153a5ebd0..9567416192c6 100644 --- a/packages/google-api-core/tests/unit/test_grpc_helpers.py +++ b/packages/google-api-core/tests/unit/test_grpc_helpers.py @@ -1033,3 +1033,140 @@ def test_apply_channel_interceptors_invalid_type_raises(): mock_base_channel = mock.Mock(name="base_channel") with pytest.raises(TypeError, match="Expected ClientInterceptor or Callable"): grpc_helpers.apply_channel_interceptors(mock_base_channel, [12345]) + + +@pytest.mark.parametrize( + "target,attempt_interconnect,expected_target", + [ + ( + "storage-direct.googleapis.com:443", + True, + "google-c2p:///storage-direct.googleapis.com?force-xds", + ), + ( + "dns:///storage-direct.googleapis.com:443", + True, + "google-c2p:///storage-direct.googleapis.com?force-xds", + ), + ( + "google-c2p:///storage-direct.googleapis.com", + True, + "google-c2p:///storage-direct.googleapis.com?force-xds", + ), + ( + "google-c2p:///storage-direct.googleapis.com?force-xds", + True, + "google-c2p:///storage-direct.googleapis.com?force-xds", + ), + ( + "google-c2p:///storage-direct.googleapis.com?force-xds", + False, + "google-c2p:///storage-direct.googleapis.com?force-xds", + ), + ( + "storage-direct.googleapis.com:443?foo=bar", + True, + "google-c2p:///storage-direct.googleapis.com?foo=bar&force-xds", + ), + ], +) +def test__modify_target_for_direct_path_interconnect( + target, attempt_interconnect, expected_target +): + actual = grpc_helpers._modify_target_for_direct_path( + target, + attempt_direct_path_xds_over_interconnect=attempt_interconnect, + ) + assert actual == expected_target + + +@mock.patch("grpc.ssl_channel_credentials") +@mock.patch("grpc.composite_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.secure_channel") +def test_create_channel_attempt_direct_path_xds_over_interconnect_default_ssl( + grpc_secure_channel, + google_auth_default, + composite_creds_call, + ssl_creds_call, +): + composite_creds = composite_creds_call.return_value + default_ssl_creds = ssl_creds_call.return_value + + channel = grpc_helpers.create_channel( + "storage-direct.googleapis.com:443", + attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=True, + ) + + assert channel is grpc_secure_channel.return_value + ssl_creds_call.assert_called_once_with() + assert composite_creds_call.call_args.args[0] is default_ssl_creds + grpc_secure_channel.assert_called_once_with( + "google-c2p:///storage-direct.googleapis.com?force-xds", + composite_creds, + compression=None, + ) + + +@mock.patch("grpc.composite_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.secure_channel") +def test_create_channel_attempt_direct_path_xds_over_interconnect_custom_ssl( + grpc_secure_channel, + google_auth_default, + composite_creds_call, +): + composite_creds = composite_creds_call.return_value + custom_ssl_creds = mock.sentinel.custom_ssl_creds + + channel = grpc_helpers.create_channel( + "storage-direct.googleapis.com:443", + ssl_credentials=custom_ssl_creds, + attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=True, + ) + + assert channel is grpc_secure_channel.return_value + assert composite_creds_call.call_args.args[0] is custom_ssl_creds + grpc_secure_channel.assert_called_once_with( + "google-c2p:///storage-direct.googleapis.com?force-xds", + composite_creds, + compression=None, + ) + + +@mock.patch("grpc.compute_engine_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.secure_channel") +def test_create_channel_interconnect_fallback_rewrites_direct_host( + grpc_secure_channel, + google_auth_default, + composite_creds_call, +): + composite_creds = composite_creds_call.return_value + + channel = grpc_helpers.create_channel( + "storage-direct.googleapis.com:443", + attempt_direct_path=False, + attempt_direct_path_xds_over_interconnect=False, + ) + + assert channel is grpc_secure_channel.return_value + grpc_secure_channel.assert_called_once_with( + "storage.googleapis.com:443", + composite_creds, + compression=None, + ) diff --git a/packages/google-cloud-storage/google/cloud/storage/asyncio/async_grpc_client.py b/packages/google-cloud-storage/google/cloud/storage/asyncio/async_grpc_client.py index 18ee348401c7..e210fb5fc195 100644 --- a/packages/google-cloud-storage/google/cloud/storage/asyncio/async_grpc_client.py +++ b/packages/google-cloud-storage/google/cloud/storage/asyncio/async_grpc_client.py @@ -22,8 +22,11 @@ DEFAULT_CLIENT_INFO, ) from google.cloud.storage import __version__ - -_DEFAULT_HOST = "storage.googleapis.com" +from google.cloud.storage.grpc_client import ( + _DEFAULT_HOST, + _resolve_direct_path_interconnect, + _rewrite_host_for_interconnect, +) def _validate_metadata(metadata): @@ -63,6 +66,13 @@ class AsyncGrpcClient: :param attempt_direct_path: (Optional) Whether to attempt to use DirectPath for gRPC connections. Defaults to ``True``. + + :type attempt_direct_path_xds_over_interconnect: bool + :param attempt_direct_path_xds_over_interconnect: + (Optional) Whether to attempt DirectPath over Cloud Interconnect + using xDS and standard TLS. Can also be overridden via the + ``GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT`` environment + variable (``"true"`` or ``"false"``). Defaults to ``False``. """ def __init__( @@ -72,6 +82,7 @@ def __init__( client_options=None, *, attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=False, ): if isinstance(credentials, auth_credentials.AnonymousCredentials): if client_options is None or client_options.api_endpoint is None: @@ -92,11 +103,16 @@ def __init__( if agent_version not in client_info.user_agent: client_info.user_agent += f" {agent_version} " + use_dp_interconnect = _resolve_direct_path_interconnect( + attempt_direct_path_xds_over_interconnect + ) + self._grpc_client = self._create_async_grpc_client( credentials=credentials, client_info=client_info, client_options=client_options, attempt_direct_path=attempt_direct_path, + attempt_direct_path_xds_over_interconnect=use_dp_interconnect, ) def _create_anonymous_client(self, client_options, credentials): @@ -120,6 +136,7 @@ def _create_async_grpc_client( client_info=None, client_options=None, attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=False, ): transport_cls = storage_v2.StorageAsyncClient.get_transport_class( "grpc_asyncio" @@ -129,17 +146,33 @@ def _create_async_grpc_client( host = _DEFAULT_HOST quota_project_id = None - if client_options: + if isinstance(client_options, dict): + host = client_options.get("api_endpoint") or _DEFAULT_HOST + quota_project_id = client_options.get("quota_project_id") + elif client_options: host = getattr(client_options, "api_endpoint", None) or _DEFAULT_HOST quota_project_id = getattr(client_options, "quota_project_id", None) - channel = transport_cls.create_channel( - host=host, - quota_project_id=quota_project_id, - attempt_direct_path=attempt_direct_path, - credentials=credentials, - options=(("grpc.primary_user_agent", primary_user_agent),), - ) + if attempt_direct_path_xds_over_interconnect: + host = _rewrite_host_for_interconnect(host) + channel = transport_cls.create_channel( + host=host, + quota_project_id=quota_project_id, + attempt_direct_path=bool( + attempt_direct_path or attempt_direct_path_xds_over_interconnect + ), + attempt_direct_path_xds_over_interconnect=True, + credentials=credentials, + options=(("grpc.primary_user_agent", primary_user_agent),), + ) + else: + channel = transport_cls.create_channel( + host=host, + quota_project_id=quota_project_id, + attempt_direct_path=attempt_direct_path, + credentials=credentials, + options=(("grpc.primary_user_agent", primary_user_agent),), + ) transport = transport_cls(channel=channel) return storage_v2.StorageAsyncClient( diff --git a/packages/google-cloud-storage/google/cloud/storage/grpc_client.py b/packages/google-cloud-storage/google/cloud/storage/grpc_client.py index 4fdb23e702b6..f1d7474c68d0 100644 --- a/packages/google-cloud-storage/google/cloud/storage/grpc_client.py +++ b/packages/google-cloud-storage/google/cloud/storage/grpc_client.py @@ -14,11 +14,43 @@ """A client for interacting with Google Cloud Storage using the gRPC API.""" +import os + from google.cloud.client import ClientWithProject from google.cloud import _storage_v2 as storage_v2 _marker = object() +_DEFAULT_HOST = "storage.googleapis.com" +_DEFAULT_HOST_DIRECT_PATH = "storage-direct.googleapis.com" +_DIRECT_PATH_INTERCONNECT_ENV = "GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT" + + +def _resolve_direct_path_interconnect(option_value: bool) -> bool: + """Resolves DirectPath over Interconnect flag from env var or parameter.""" + env_val = os.environ.get(_DIRECT_PATH_INTERCONNECT_ENV) + if env_val == "true": + return True + if env_val == "false": + return False + return bool(option_value) + + +def _rewrite_host_for_interconnect( + host: str, + old_host: str = _DEFAULT_HOST, + new_host: str = _DEFAULT_HOST_DIRECT_PATH, +) -> str: + """Rewrites the default GCS endpoint to the DirectPath over Interconnect host.""" + if not host: + return new_host + for prefix in ("", "https://", "http://", "dns:///"): + candidate = f"{prefix}{old_host}" + if host.startswith(candidate): + rest = host[len(candidate) :] + if not rest or rest[0] in (":", "/", "?", "#"): + return f"{prefix}{new_host}{rest}" + return host class GrpcClient(ClientWithProject): @@ -58,6 +90,13 @@ 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 attempt_direct_path_xds_over_interconnect: bool + :param attempt_direct_path_xds_over_interconnect: + (Optional) Whether to attempt DirectPath over Cloud Interconnect + using xDS and standard TLS. Can also be overridden via the + ``GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT`` environment + variable (``"true"`` or ``"false"``). Defaults to ``False``. """ def __init__( @@ -69,6 +108,7 @@ def __init__( *, api_key=None, attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=False, ): super(GrpcClient, self).__init__(project=project, credentials=credentials) @@ -80,11 +120,16 @@ def __init__( elif api_key: client_options.api_key = api_key + use_dp_interconnect = _resolve_direct_path_interconnect( + attempt_direct_path_xds_over_interconnect + ) + self._grpc_client = self._create_gapic_client( credentials=credentials, client_info=client_info, client_options=client_options, attempt_direct_path=attempt_direct_path, + attempt_direct_path_xds_over_interconnect=use_dp_interconnect, ) def _create_gapic_client( @@ -93,11 +138,29 @@ def _create_gapic_client( client_info=None, client_options=None, attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=False, ): """Creates and configures the low-level GAPIC `storage_v2` client.""" transport_cls = storage_v2.StorageClient.get_transport_class("grpc") - channel = transport_cls.create_channel(attempt_direct_path=attempt_direct_path) + if attempt_direct_path_xds_over_interconnect: + host = _DEFAULT_HOST + if isinstance(client_options, dict): + host = client_options.get("api_endpoint") or _DEFAULT_HOST + elif client_options is not None: + host = getattr(client_options, "api_endpoint", None) or _DEFAULT_HOST + host = _rewrite_host_for_interconnect(host) + channel = transport_cls.create_channel( + host=host, + attempt_direct_path=bool( + attempt_direct_path or attempt_direct_path_xds_over_interconnect + ), + attempt_direct_path_xds_over_interconnect=True, + ) + else: + channel = transport_cls.create_channel( + attempt_direct_path=attempt_direct_path + ) transport = transport_cls(credentials=credentials, channel=channel) diff --git a/packages/google-cloud-storage/tests/unit/asyncio/test_async_grpc_client.py b/packages/google-cloud-storage/tests/unit/asyncio/test_async_grpc_client.py index 24168dc137d4..961eb8b1529e 100644 --- a/packages/google-cloud-storage/tests/unit/asyncio/test_async_grpc_client.py +++ b/packages/google-cloud-storage/tests/unit/asyncio/test_async_grpc_client.py @@ -431,3 +431,62 @@ async def test_get_object_invalid_metadata(self, invalid_metadata, expected_exc) client = async_grpc_client.AsyncGrpcClient(credentials=_make_credentials()) with pytest.raises(expected_exc): await client.get_object("bucket", "object", metadata=invalid_metadata) + + @mock.patch("google.cloud._storage_v2.StorageAsyncClient") + def test_async_grpc_client_attempt_direct_path_xds_over_interconnect( + self, mock_async_storage_client + ): + mock_transport_cls = mock.MagicMock() + mock_async_storage_client.get_transport_class.return_value = mock_transport_cls + mock_creds = _make_credentials() + + async_grpc_client.AsyncGrpcClient( + credentials=mock_creds, + attempt_direct_path_xds_over_interconnect=True, + ) + + kwargs = mock_async_storage_client.call_args.kwargs + client_info = kwargs["client_info"] + primary_user_agent = client_info.to_user_agent() + expected_options = (("grpc.primary_user_agent", primary_user_agent),) + + mock_transport_cls.create_channel.assert_called_once_with( + host="storage-direct.googleapis.com", + quota_project_id=None, + attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=True, + credentials=mock_creds, + options=expected_options, + ) + + @mock.patch("google.cloud._storage_v2.StorageAsyncClient") + def test_async_grpc_client_interconnect_env_override( + self, mock_async_storage_client + ): + mock_transport_cls = mock.MagicMock() + mock_async_storage_client.get_transport_class.return_value = mock_transport_cls + mock_creds = _make_credentials() + + with mock.patch.dict( + "os.environ", + {"GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT": "true"}, + ): + async_grpc_client.AsyncGrpcClient( + credentials=mock_creds, + attempt_direct_path=False, + attempt_direct_path_xds_over_interconnect=False, + ) + + kwargs = mock_async_storage_client.call_args.kwargs + client_info = kwargs["client_info"] + primary_user_agent = client_info.to_user_agent() + expected_options = (("grpc.primary_user_agent", primary_user_agent),) + + mock_transport_cls.create_channel.assert_called_once_with( + host="storage-direct.googleapis.com", + quota_project_id=None, + attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=True, + credentials=mock_creds, + options=expected_options, + ) diff --git a/packages/google-cloud-storage/tests/unit/test_grpc_client.py b/packages/google-cloud-storage/tests/unit/test_grpc_client.py index 9baa15c0ad5f..9ca14fc7b0d5 100644 --- a/packages/google-cloud-storage/tests/unit/test_grpc_client.py +++ b/packages/google-cloud-storage/tests/unit/test_grpc_client.py @@ -197,3 +197,100 @@ 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_grpc_client_attempt_direct_path_xds_over_interconnect( + self, mock_storage_client, mock_base_client + ): + mock_transport_cls = mock.MagicMock() + mock_storage_client.get_transport_class.return_value = mock_transport_cls + mock_creds = _make_credentials() + mock_base_client.return_value._credentials = mock_creds + + grpc_client.GrpcClient( + project="test-project", + credentials=mock_creds, + attempt_direct_path_xds_over_interconnect=True, + ) + + mock_transport_cls.create_channel.assert_called_once_with( + host="storage-direct.googleapis.com", + attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=True, + ) + + @mock.patch("google.cloud.storage.grpc_client.ClientWithProject") + @mock.patch("google.cloud._storage_v2.StorageClient") + def test_grpc_client_attempt_direct_path_xds_over_interconnect_env_override( + self, mock_storage_client, mock_base_client + ): + mock_transport_cls = mock.MagicMock() + mock_storage_client.get_transport_class.return_value = mock_transport_cls + mock_creds = _make_credentials() + mock_base_client.return_value._credentials = mock_creds + + with mock.patch.dict( + "os.environ", + {"GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT": "true"}, + ): + grpc_client.GrpcClient( + project="test-project", + credentials=mock_creds, + attempt_direct_path_xds_over_interconnect=False, + ) + + mock_transport_cls.create_channel.assert_called_once_with( + host="storage-direct.googleapis.com", + attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=True, + ) + + mock_transport_cls.create_channel.reset_mock() + with mock.patch.dict( + "os.environ", + {"GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT": "false"}, + ): + grpc_client.GrpcClient( + project="test-project", + credentials=mock_creds, + attempt_direct_path_xds_over_interconnect=True, + ) + + mock_transport_cls.create_channel.assert_called_once_with( + attempt_direct_path=True + ) + + def test_rewrite_host_for_interconnect_delimiters(self): + cases = [ + ("storage.googleapis.com", "storage-direct.googleapis.com"), + ("storage.googleapis.com:443", "storage-direct.googleapis.com:443"), + ("storage.googleapis.com/path", "storage-direct.googleapis.com/path"), + ( + "storage.googleapis.com?query=val", + "storage-direct.googleapis.com?query=val", + ), + ("storage.googleapis.com#section", "storage-direct.googleapis.com#section"), + ("https://storage.googleapis.com", "https://storage-direct.googleapis.com"), + ( + "https://storage.googleapis.com:443?q=1", + "https://storage-direct.googleapis.com:443?q=1", + ), + ( + "dns:///storage.googleapis.com:443", + "dns:///storage-direct.googleapis.com:443", + ), + ( + "storage.googleapis.com.evil.com", + "storage.googleapis.com.evil.com", + ), + ( + "google-c2p:///storage-direct.googleapis.com", + "google-c2p:///storage-direct.googleapis.com", + ), + ] + for original_host, expected_host in cases: + self.assertEqual( + grpc_client._rewrite_host_for_interconnect(original_host), + expected_host, + ) From 2e058acb87d21edbd17bb7e92b7c271294f5a707 Mon Sep 17 00:00:00 2001 From: Nidhi Nandwani Date: Sun, 20 Sep 2026 13:38:26 +0000 Subject: [PATCH 2/4] test(api_core): add unit tests for DirectPath over Interconnect coverage Ensure 100% test coverage for _resolve_direct_path_interconnect and _create_composite_credentials in google.api_core.grpc_helpers. [Generated-by: AI] --- .../tests/unit/test_grpc_helpers.py | 48 +++++++++++++++++++ 1 file changed, 48 insertions(+) diff --git a/packages/google-api-core/tests/unit/test_grpc_helpers.py b/packages/google-api-core/tests/unit/test_grpc_helpers.py index 9567416192c6..ff56bc296ca5 100644 --- a/packages/google-api-core/tests/unit/test_grpc_helpers.py +++ b/packages/google-api-core/tests/unit/test_grpc_helpers.py @@ -1035,6 +1035,54 @@ def test_apply_channel_interceptors_invalid_type_raises(): grpc_helpers.apply_channel_interceptors(mock_base_channel, [12345]) +@pytest.mark.parametrize( + "env_value,attempt_interconnect,expected_result", + [ + ("true", False, True), + ("true", None, True), + ("false", True, False), + ("false", None, False), + (None, True, True), + (None, False, False), + (None, None, False), + ], +) +def test__resolve_direct_path_interconnect( + monkeypatch, env_value, attempt_interconnect, expected_result +): + if env_value is None: + monkeypatch.delenv(grpc_helpers._DIRECT_PATH_INTERCONNECT_ENV, raising=False) + else: + monkeypatch.setenv(grpc_helpers._DIRECT_PATH_INTERCONNECT_ENV, env_value) + assert ( + grpc_helpers._resolve_direct_path_interconnect(attempt_interconnect) + is expected_result + ) + + +@mock.patch("grpc.compute_engine_channel_credentials") +@mock.patch("google.auth.transport.requests.Request", autospec=True) +@mock.patch("google.auth.transport.grpc.AuthMetadataPlugin") +def test__composite_credentials_auth_metadata_plugin_type_error_fallback( + auth_metadata_plugin, request, compute_engine_creds +): + fallback_plugin = mock.sentinel.fallback_plugin + auth_metadata_plugin.side_effect = [ + TypeError("unexpected keyword"), + fallback_plugin, + ] + credentials = mock.create_autospec( + google.auth.credentials.Credentials, instance=True + ) + credentials.requires_scopes = False + + grpc_helpers._create_composite_credentials( + credentials, default_host="storage.googleapis.com" + ) + + assert auth_metadata_plugin.call_count == 2 + + @pytest.mark.parametrize( "target,attempt_interconnect,expected_target", [ From 04d5cb8a46fff90bb21015863262f2a56ef78742 Mon Sep 17 00:00:00 2001 From: Nidhi Nandwani Date: Mon, 21 Sep 2026 11:49:18 +0000 Subject: [PATCH 3/4] fix(storage,api_core): resolve PR review comments and add DirectPath authority override - Normalize casing and whitespace in _resolve_direct_path_interconnect and raise ValueError on invalid values - Inject grpc.ssl_target_name_override authority override for -direct.googleapis.com endpoints when DirectPath over Interconnect is enabled - Restrict CloudPath fallback host replacement to -direct.googleapis.com -> .googleapis.com - Strip dns:///, https://, and http:// scheme prefixes and path components in _modify_target_for_direct_path - Forward credentials and quota_project_id in GrpcClient._create_gapic_client when attempt_direct_path_xds_over_interconnect is enabled [Generated-by: AI] --- .../google/api_core/grpc_helpers.py | 61 ++++++++++++++----- .../google/api_core/grpc_helpers_async.py | 17 +++++- .../tests/asyncio/test_grpc_helpers_async.py | 41 +++++++++++++ .../tests/unit/test_grpc_helpers.py | 58 +++++++++++++++++- .../google/cloud/storage/grpc_client.py | 31 +++++++--- .../unit/asyncio/test_async_grpc_client.py | 7 +++ .../tests/unit/test_grpc_client.py | 41 ++++++++++++- 7 files changed, 227 insertions(+), 29 deletions(-) diff --git a/packages/google-api-core/google/api_core/grpc_helpers.py b/packages/google-api-core/google/api_core/grpc_helpers.py index 27fca814425f..0e52cc7b158b 100644 --- a/packages/google-api-core/google/api_core/grpc_helpers.py +++ b/packages/google-api-core/google/api_core/grpc_helpers.py @@ -47,13 +47,32 @@ def _resolve_direct_path_interconnect( ) -> bool: """Resolves whether DirectPath over Interconnect is enabled.""" env_val = os.environ.get(_DIRECT_PATH_INTERCONNECT_ENV) - if env_val == "true": - return True - if env_val == "false": - return False + if env_val is not None: + env_val_clean = env_val.strip().lower() + if env_val_clean == "true": + return True + elif env_val_clean == "false": + return False + else: + raise ValueError( + f"Invalid value for {_DIRECT_PATH_INTERCONNECT_ENV}: {env_val}" + ) return bool(attempt_direct_path_xds_over_interconnect) +def _extract_direct_path_authority(target: str) -> Optional[str]: + """Extracts the canonical TLS/HTTP2 authority for a ``-direct.googleapis.com`` target.""" + clean_host = target + for prefix in ("google-c2p:///", "dns:///", "https://", "http://"): + if clean_host.startswith(prefix): + clean_host = clean_host[len(prefix) :] + break + clean_host = clean_host.split("?", 1)[0].split("/", 1)[0].split(":", 1)[0] + if "-direct.googleapis.com" in clean_host: + return clean_host.replace("-direct.googleapis.com", ".googleapis.com", 1) + return None + + # The list of gRPC Callable interfaces that return iterators. _STREAM_WRAP_CLASSES = (grpc.UnaryStreamMultiCallable, grpc.StreamStreamMultiCallable) @@ -431,13 +450,26 @@ def create_channel( default_host=default_host, ) + if use_dp_interconnect: + authority = _extract_direct_path_authority(target) + if authority: + existing_options = tuple(kwargs.get("options") or ()) + option_keys = {opt[0] for opt in existing_options} + if ( + "grpc.ssl_target_name_override" not in option_keys + and "grpc.default_authority" not in option_keys + ): + kwargs["options"] = existing_options + ( + ("grpc.ssl_target_name_override", authority), + ) + if attempt_direct_path or use_dp_interconnect: target = _modify_target_for_direct_path( target, attempt_direct_path_xds_over_interconnect=use_dp_interconnect, ) - elif "-direct." in target and not target.startswith("google-c2p:///"): - target = target.replace("-direct.", ".") + elif "-direct.googleapis.com" in target and not target.startswith("google-c2p:///"): + target = target.replace("-direct.googleapis.com", ".googleapis.com") return grpc.secure_channel( target, composite_credentials, compression=compression, **kwargs @@ -464,23 +496,24 @@ def _modify_target_for_direct_path( original target may already denote Direct Path. """ - # A DNS prefix may be included with the target to indicate the endpoint is living in the Internet, - # outside of Google Cloud Platform. - dns_prefix = "dns:///" - # Remove "dns:///" if `attempt_direct_path` is set to True as - # the Direct Path prefix `google-c2p:///` will be used instead. - target = target.replace(dns_prefix, "") + # Strip standard URI scheme prefixes ("dns:///", "https://", "http://") if + # `attempt_direct_path` is enabled, as the Direct Path prefix `google-c2p:///` + # will be used instead. + for scheme_prefix in ("dns:///", "https://", "http://"): + if target.startswith(scheme_prefix): + target = target[len(scheme_prefix) :] + break direct_path_separator = ":///" if direct_path_separator not in target: if "?" in target: host_part, query_part = target.split("?", 1) - target_without_port = host_part.split(":")[0] + target_without_port = host_part.split("/")[0].split(":")[0] target = ( f"google-c2p{direct_path_separator}{target_without_port}?{query_part}" ) else: - target_without_port = target.split(":")[0] + target_without_port = target.split("/")[0].split(":")[0] # Modify the target to use Direct Path by adding the `google-c2p:///` prefix target = f"google-c2p{direct_path_separator}{target_without_port}" diff --git a/packages/google-api-core/google/api_core/grpc_helpers_async.py b/packages/google-api-core/google/api_core/grpc_helpers_async.py index 13940885b4c7..4323fd867037 100644 --- a/packages/google-api-core/google/api_core/grpc_helpers_async.py +++ b/packages/google-api-core/google/api_core/grpc_helpers_async.py @@ -313,13 +313,26 @@ def create_channel( default_host=default_host, ) + if use_dp_interconnect: + authority = grpc_helpers._extract_direct_path_authority(target) + if authority: + existing_options = tuple(kwargs.get("options") or ()) + option_keys = {opt[0] for opt in existing_options} + if ( + "grpc.ssl_target_name_override" not in option_keys + and "grpc.default_authority" not in option_keys + ): + kwargs["options"] = existing_options + ( + ("grpc.ssl_target_name_override", authority), + ) + if attempt_direct_path or use_dp_interconnect: target = grpc_helpers._modify_target_for_direct_path( target, attempt_direct_path_xds_over_interconnect=use_dp_interconnect, ) - elif "-direct." in target and not target.startswith("google-c2p:///"): - target = target.replace("-direct.", ".") + elif "-direct.googleapis.com" in target and not target.startswith("google-c2p:///"): + target = target.replace("-direct.googleapis.com", ".googleapis.com") return aio.secure_channel( target, composite_credentials, compression=compression, **kwargs diff --git a/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py b/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py index bab77fe765c4..8dc0f5c6d127 100644 --- a/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py +++ b/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py @@ -775,6 +775,7 @@ def test_create_channel_attempt_direct_path_xds_over_interconnect_default_ssl( "google-c2p:///storage-direct.googleapis.com?force-xds", composite_creds, compression=None, + options=(("grpc.ssl_target_name_override", "storage.googleapis.com"),), ) @@ -806,6 +807,7 @@ def test_create_channel_attempt_direct_path_xds_over_interconnect_custom_ssl( "google-c2p:///storage-direct.googleapis.com?force-xds", composite_creds, compression=None, + options=(("grpc.ssl_target_name_override", "storage.googleapis.com"),), ) @@ -835,3 +837,42 @@ def test_create_channel_interconnect_fallback( composite_creds, compression=None, ) + + +@mock.patch("grpc.compute_engine_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.aio.secure_channel") +def test_create_channel_interconnect_fallback_preserves_non_google_direct_host( + grpc_secure_channel, + google_auth_default, + composite_creds_call, +): + composite_creds = composite_creds_call.return_value + + channel = grpc_helpers_async.create_channel( + "my-direct.example.com:443", + attempt_direct_path=False, + attempt_direct_path_xds_over_interconnect=False, + ) + + assert channel is grpc_secure_channel.return_value + grpc_secure_channel.assert_called_once_with( + "my-direct.example.com:443", + composite_creds, + compression=None, + ) + + +def test_create_channel_async_invalid_env_var_raises(monkeypatch): + monkeypatch.setenv( + "GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT", "invalid" + ) + with pytest.raises( + ValueError, + match="Invalid value for GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT", + ): + grpc_helpers_async.create_channel("storage-direct.googleapis.com:443") diff --git a/packages/google-api-core/tests/unit/test_grpc_helpers.py b/packages/google-api-core/tests/unit/test_grpc_helpers.py index ff56bc296ca5..57846b384270 100644 --- a/packages/google-api-core/tests/unit/test_grpc_helpers.py +++ b/packages/google-api-core/tests/unit/test_grpc_helpers.py @@ -1039,9 +1039,11 @@ def test_apply_channel_interceptors_invalid_type_raises(): "env_value,attempt_interconnect,expected_result", [ ("true", False, True), - ("true", None, True), + ("TRUE", False, True), + (" True ", None, True), ("false", True, False), - ("false", None, False), + ("FALSE", True, False), + (" False ", None, False), (None, True, True), (None, False, False), (None, None, False), @@ -1060,6 +1062,18 @@ def test__resolve_direct_path_interconnect( ) +@pytest.mark.parametrize("invalid_env_value", ["invalid", "1", "0", " ", ""]) +def test__resolve_direct_path_interconnect_invalid_raises( + monkeypatch, invalid_env_value +): + monkeypatch.setenv(grpc_helpers._DIRECT_PATH_INTERCONNECT_ENV, invalid_env_value) + with pytest.raises( + ValueError, + match=f"Invalid value for {grpc_helpers._DIRECT_PATH_INTERCONNECT_ENV}", + ): + grpc_helpers._resolve_direct_path_interconnect(False) + + @mock.patch("grpc.compute_engine_channel_credentials") @mock.patch("google.auth.transport.requests.Request", autospec=True) @mock.patch("google.auth.transport.grpc.AuthMetadataPlugin") @@ -1096,6 +1110,16 @@ def test__composite_credentials_auth_metadata_plugin_type_error_fallback( True, "google-c2p:///storage-direct.googleapis.com?force-xds", ), + ( + "https://storage-direct.googleapis.com:443", + True, + "google-c2p:///storage-direct.googleapis.com?force-xds", + ), + ( + "http://storage-direct.googleapis.com:443/v1", + True, + "google-c2p:///storage-direct.googleapis.com?force-xds", + ), ( "google-c2p:///storage-direct.googleapis.com", True, @@ -1158,6 +1182,7 @@ def test_create_channel_attempt_direct_path_xds_over_interconnect_default_ssl( "google-c2p:///storage-direct.googleapis.com?force-xds", composite_creds, compression=None, + options=(("grpc.ssl_target_name_override", "storage.googleapis.com"),), ) @@ -1189,6 +1214,7 @@ def test_create_channel_attempt_direct_path_xds_over_interconnect_custom_ssl( "google-c2p:///storage-direct.googleapis.com?force-xds", composite_creds, compression=None, + options=(("grpc.ssl_target_name_override", "storage.googleapis.com"),), ) @@ -1218,3 +1244,31 @@ def test_create_channel_interconnect_fallback_rewrites_direct_host( composite_creds, compression=None, ) + + +@mock.patch("grpc.compute_engine_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.secure_channel") +def test_create_channel_interconnect_fallback_preserves_non_google_direct_host( + grpc_secure_channel, + google_auth_default, + composite_creds_call, +): + composite_creds = composite_creds_call.return_value + + channel = grpc_helpers.create_channel( + "my-direct.example.com:443", + attempt_direct_path=False, + attempt_direct_path_xds_over_interconnect=False, + ) + + assert channel is grpc_secure_channel.return_value + grpc_secure_channel.assert_called_once_with( + "my-direct.example.com:443", + composite_creds, + compression=None, + ) diff --git a/packages/google-cloud-storage/google/cloud/storage/grpc_client.py b/packages/google-cloud-storage/google/cloud/storage/grpc_client.py index f1d7474c68d0..01415c9640f4 100644 --- a/packages/google-cloud-storage/google/cloud/storage/grpc_client.py +++ b/packages/google-cloud-storage/google/cloud/storage/grpc_client.py @@ -29,10 +29,16 @@ def _resolve_direct_path_interconnect(option_value: bool) -> bool: """Resolves DirectPath over Interconnect flag from env var or parameter.""" env_val = os.environ.get(_DIRECT_PATH_INTERCONNECT_ENV) - if env_val == "true": - return True - if env_val == "false": - return False + if env_val is not None: + env_val_clean = env_val.strip().lower() + if env_val_clean == "true": + return True + elif env_val_clean == "false": + return False + else: + raise ValueError( + f"Invalid value for {_DIRECT_PATH_INTERCONNECT_ENV}: {env_val}" + ) return bool(option_value) @@ -145,18 +151,25 @@ def _create_gapic_client( if attempt_direct_path_xds_over_interconnect: host = _DEFAULT_HOST + quota_project_id = None if isinstance(client_options, dict): host = client_options.get("api_endpoint") or _DEFAULT_HOST + quota_project_id = client_options.get("quota_project_id") elif client_options is not None: host = getattr(client_options, "api_endpoint", None) or _DEFAULT_HOST + quota_project_id = getattr(client_options, "quota_project_id", None) host = _rewrite_host_for_interconnect(host) - channel = transport_cls.create_channel( - host=host, - attempt_direct_path=bool( + channel_kwargs = { + "host": host, + "credentials": credentials, + "attempt_direct_path": bool( attempt_direct_path or attempt_direct_path_xds_over_interconnect ), - attempt_direct_path_xds_over_interconnect=True, - ) + "attempt_direct_path_xds_over_interconnect": True, + } + if quota_project_id is not None: + channel_kwargs["quota_project_id"] = quota_project_id + channel = transport_cls.create_channel(**channel_kwargs) else: channel = transport_cls.create_channel( attempt_direct_path=attempt_direct_path diff --git a/packages/google-cloud-storage/tests/unit/asyncio/test_async_grpc_client.py b/packages/google-cloud-storage/tests/unit/asyncio/test_async_grpc_client.py index 961eb8b1529e..db2c3109b5b9 100644 --- a/packages/google-cloud-storage/tests/unit/asyncio/test_async_grpc_client.py +++ b/packages/google-cloud-storage/tests/unit/asyncio/test_async_grpc_client.py @@ -490,3 +490,10 @@ def test_async_grpc_client_interconnect_env_override( credentials=mock_creds, options=expected_options, ) + + with mock.patch.dict( + "os.environ", + {"GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT": "invalid"}, + ): + with pytest.raises(ValueError): + async_grpc_client.AsyncGrpcClient(credentials=mock_creds) diff --git a/packages/google-cloud-storage/tests/unit/test_grpc_client.py b/packages/google-cloud-storage/tests/unit/test_grpc_client.py index 9ca14fc7b0d5..54b82a77d61b 100644 --- a/packages/google-cloud-storage/tests/unit/test_grpc_client.py +++ b/packages/google-cloud-storage/tests/unit/test_grpc_client.py @@ -216,6 +216,32 @@ def test_grpc_client_attempt_direct_path_xds_over_interconnect( mock_transport_cls.create_channel.assert_called_once_with( host="storage-direct.googleapis.com", + credentials=mock_creds, + attempt_direct_path=True, + attempt_direct_path_xds_over_interconnect=True, + ) + + @mock.patch("google.cloud.storage.grpc_client.ClientWithProject") + @mock.patch("google.cloud._storage_v2.StorageClient") + def test_grpc_client_attempt_direct_path_xds_over_interconnect_with_quota_project( + self, mock_storage_client, mock_base_client + ): + mock_transport_cls = mock.MagicMock() + mock_storage_client.get_transport_class.return_value = mock_transport_cls + mock_creds = _make_credentials() + mock_base_client.return_value._credentials = mock_creds + + grpc_client.GrpcClient( + project="test-project", + credentials=mock_creds, + client_options={"quota_project_id": "my-quota-proj"}, + attempt_direct_path_xds_over_interconnect=True, + ) + + mock_transport_cls.create_channel.assert_called_once_with( + host="storage-direct.googleapis.com", + credentials=mock_creds, + quota_project_id="my-quota-proj", attempt_direct_path=True, attempt_direct_path_xds_over_interconnect=True, ) @@ -232,7 +258,7 @@ def test_grpc_client_attempt_direct_path_xds_over_interconnect_env_override( with mock.patch.dict( "os.environ", - {"GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT": "true"}, + {"GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT": " TRUE "}, ): grpc_client.GrpcClient( project="test-project", @@ -242,6 +268,7 @@ def test_grpc_client_attempt_direct_path_xds_over_interconnect_env_override( mock_transport_cls.create_channel.assert_called_once_with( host="storage-direct.googleapis.com", + credentials=mock_creds, attempt_direct_path=True, attempt_direct_path_xds_over_interconnect=True, ) @@ -249,7 +276,7 @@ def test_grpc_client_attempt_direct_path_xds_over_interconnect_env_override( mock_transport_cls.create_channel.reset_mock() with mock.patch.dict( "os.environ", - {"GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT": "false"}, + {"GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT": " FALSE "}, ): grpc_client.GrpcClient( project="test-project", @@ -261,6 +288,16 @@ def test_grpc_client_attempt_direct_path_xds_over_interconnect_env_override( attempt_direct_path=True ) + with mock.patch.dict( + "os.environ", + {"GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT": "invalid"}, + ): + with self.assertRaises(ValueError): + grpc_client.GrpcClient( + project="test-project", + credentials=mock_creds, + ) + def test_rewrite_host_for_interconnect_delimiters(self): cases = [ ("storage.googleapis.com", "storage-direct.googleapis.com"), From d8a2d4eedf4ea067c0bcbfe9a2dc273ab4219606 Mon Sep 17 00:00:00 2001 From: Nidhi Nandwani Date: Mon, 21 Sep 2026 12:05:07 +0000 Subject: [PATCH 4/4] test(api_core): cover _extract_direct_path_authority and authority override branches Add unit tests for _extract_direct_path_authority and authority override branches in grpc_helpers and grpc_helpers_async to restore 100% coverage. [Generated-by: AI] --- .../tests/asyncio/test_grpc_helpers_async.py | 42 +++++++++++++ .../tests/unit/test_grpc_helpers.py | 59 +++++++++++++++++++ 2 files changed, 101 insertions(+) diff --git a/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py b/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py index 8dc0f5c6d127..9c5723843c49 100644 --- a/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py +++ b/packages/google-api-core/tests/asyncio/test_grpc_helpers_async.py @@ -876,3 +876,45 @@ def test_create_channel_async_invalid_env_var_raises(monkeypatch): match="Invalid value for GOOGLE_CLOUD_ENABLE_DIRECT_PATH_XDS_OVER_INTERCONNECT", ): grpc_helpers_async.create_channel("storage-direct.googleapis.com:443") + + +@mock.patch("grpc.ssl_channel_credentials") +@mock.patch("grpc.composite_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.aio.secure_channel") +def test_create_channel_async_interconnect_authority_branches( + grpc_secure_channel, + google_auth_default, + composite_creds_call, + ssl_creds_call, +): + composite_creds = composite_creds_call.return_value + + # Branch 318->329: authority is None (target does not contain -direct.googleapis.com) + grpc_helpers_async.create_channel( + "storage.googleapis.com:443", + attempt_direct_path_xds_over_interconnect=True, + ) + grpc_secure_channel.assert_called_with( + "google-c2p:///storage.googleapis.com?force-xds", + composite_creds, + compression=None, + ) + + # Branch 321->329: authority exists, but grpc.ssl_target_name_override is already provided + custom_options = (("grpc.ssl_target_name_override", "custom.googleapis.com"),) + grpc_helpers_async.create_channel( + "storage-direct.googleapis.com:443", + attempt_direct_path_xds_over_interconnect=True, + options=custom_options, + ) + grpc_secure_channel.assert_called_with( + "google-c2p:///storage-direct.googleapis.com?force-xds", + composite_creds, + compression=None, + options=custom_options, + ) diff --git a/packages/google-api-core/tests/unit/test_grpc_helpers.py b/packages/google-api-core/tests/unit/test_grpc_helpers.py index 57846b384270..5611f8a6346f 100644 --- a/packages/google-api-core/tests/unit/test_grpc_helpers.py +++ b/packages/google-api-core/tests/unit/test_grpc_helpers.py @@ -1272,3 +1272,62 @@ def test_create_channel_interconnect_fallback_preserves_non_google_direct_host( composite_creds, compression=None, ) + + +@pytest.mark.parametrize( + "target,expected_authority", + [ + ("storage-direct.googleapis.com:443", "storage.googleapis.com"), + ("dns:///storage-direct.googleapis.com:443", "storage.googleapis.com"), + ( + "google-c2p:///storage-direct.googleapis.com?force-xds", + "storage.googleapis.com", + ), + ("storage.googleapis.com:443", None), + ("my-direct.example.com:443", None), + ], +) +def test__extract_direct_path_authority(target, expected_authority): + assert grpc_helpers._extract_direct_path_authority(target) == expected_authority + + +@mock.patch("grpc.ssl_channel_credentials") +@mock.patch("grpc.composite_channel_credentials") +@mock.patch( + "google.auth.default", + autospec=True, + return_value=(mock.sentinel.credentials, mock.sentinel.project), +) +@mock.patch("grpc.secure_channel") +def test_create_channel_interconnect_authority_branches( + grpc_secure_channel, + google_auth_default, + composite_creds_call, + ssl_creds_call, +): + composite_creds = composite_creds_call.return_value + + # Branch 455->466: authority is None (target does not contain -direct.googleapis.com) + grpc_helpers.create_channel( + "storage.googleapis.com:443", + attempt_direct_path_xds_over_interconnect=True, + ) + grpc_secure_channel.assert_called_with( + "google-c2p:///storage.googleapis.com?force-xds", + composite_creds, + compression=None, + ) + + # Branch 458->466: authority exists, but grpc.ssl_target_name_override is already provided + custom_options = (("grpc.ssl_target_name_override", "custom.googleapis.com"),) + grpc_helpers.create_channel( + "storage-direct.googleapis.com:443", + attempt_direct_path_xds_over_interconnect=True, + options=custom_options, + ) + grpc_secure_channel.assert_called_with( + "google-c2p:///storage-direct.googleapis.com?force-xds", + composite_creds, + compression=None, + options=custom_options, + )