diff --git a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/_compat.py.j2 b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/_compat.py.j2 index 3c14eb41c04b..5380c9e26d8a 100644 --- a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/_compat.py.j2 +++ b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/_compat.py.j2 @@ -81,9 +81,19 @@ else: # pragma: NO COVER record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/_compat.py b/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/_compat.py index 5946aadc8f20..8a396c9e8e41 100755 --- a/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/_compat.py +++ b/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/_compat.py @@ -72,9 +72,19 @@ def record_error(self, exc: BaseException | None) -> None: record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/_compat.py b/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/_compat.py index 5946aadc8f20..8a396c9e8e41 100755 --- a/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/_compat.py +++ b/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/_compat.py @@ -72,9 +72,19 @@ def record_error(self, exc: BaseException | None) -> None: record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/_compat.py b/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/_compat.py index 5946aadc8f20..8a396c9e8e41 100755 --- a/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/_compat.py +++ b/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/_compat.py @@ -72,9 +72,19 @@ def record_error(self, exc: BaseException | None) -> None: record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/_compat.py b/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/_compat.py index 5946aadc8f20..8a396c9e8e41 100755 --- a/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/_compat.py +++ b/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/_compat.py @@ -72,9 +72,19 @@ def record_error(self, exc: BaseException | None) -> None: record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/_compat.py b/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/_compat.py index 5946aadc8f20..8a396c9e8e41 100755 --- a/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/_compat.py +++ b/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/_compat.py @@ -72,9 +72,19 @@ def record_error(self, exc: BaseException | None) -> None: record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/_compat.py b/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/_compat.py index 5946aadc8f20..8a396c9e8e41 100755 --- a/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/_compat.py +++ b/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/_compat.py @@ -72,9 +72,19 @@ def record_error(self, exc: BaseException | None) -> None: record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/_compat.py b/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/_compat.py index 5946aadc8f20..8a396c9e8e41 100755 --- a/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/_compat.py +++ b/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/_compat.py @@ -72,9 +72,19 @@ def record_error(self, exc: BaseException | None) -> None: record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/integration/goldens/showcase/google/showcase_v1beta1/_compat.py b/packages/gapic-generator/tests/integration/goldens/showcase/google/showcase_v1beta1/_compat.py index e9202c61ba46..f77ff8ac7590 100755 --- a/packages/gapic-generator/tests/integration/goldens/showcase/google/showcase_v1beta1/_compat.py +++ b/packages/gapic-generator/tests/integration/goldens/showcase/google/showcase_v1beta1/_compat.py @@ -77,9 +77,19 @@ def record_error(self, exc: BaseException | None) -> None: record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/_compat.py b/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/_compat.py index e9202c61ba46..f77ff8ac7590 100755 --- a/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/_compat.py +++ b/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/_compat.py @@ -77,9 +77,19 @@ def record_error(self, exc: BaseException | None) -> None: record_http_error = record_error - def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: + # mypy: google-api-core < 2.42.0 lacks trace_http_request, pragma ignores the fallback function redefinition. + def trace_http_request(*args: Any, **kwargs: Any) -> _FallbackTraceContext: # type: ignore[misc] return _FallbackTraceContext() +# OpenTelemetry gRPC auto-instrumentation suppression was introduced in +# google-api-core 2.43.0+ to prevent redundant wire spans when OpenTelemetry +# instrumentation is active alongside Google Cloud SDK tracing. +# TODO(observability): Remove once setup.py.j2 enforces google-api-core >= 2.43.0. +HAS_AUTO_INSTRUMENTATION_SUPPRESSION = ( + _observability is not None + and hasattr(_observability, "_AsyncSuppressingClientInterceptor") +) + # The `kind` parameter in gapic_v1.method_async.wrap_method was introduced in # google-api-core 2.29.0 (PR #688) alongside _DEFAULT_ASYNC_TRANSPORT_KIND to prevent # async REST callables from being wrapped with gRPC error handlers. diff --git a/packages/gapic-generator/tests/system/test_tracing.py b/packages/gapic-generator/tests/system/test_tracing.py index 7f2ca695a60b..ad733c6a3e1c 100644 --- a/packages/gapic-generator/tests/system/test_tracing.py +++ b/packages/gapic-generator/tests/system/test_tracing.py @@ -50,7 +50,10 @@ from google.api_core._feature_gating_helpers import FeatureGatingError from google.api_core.client_options import ClientOptions from google.auth import credentials as ga_credentials -from google.showcase import EchoClient +from google.showcase import EchoAsyncClient, EchoClient +from google.showcase_v1beta1._compat import ( + HAS_AUTO_INSTRUMENTATION_SUPPRESSION, +) try: from .conftest import construct_client @@ -233,3 +236,126 @@ def test_env_var_opt_in(otel_echo_client): assert len(spans) == 2 for span in spans: assert span.name == "google.showcase.v1beta1.Echo/Echo" + + +@pytest.mark.skipif( + not HAS_AUTO_INSTRUMENTATION_SUPPRESSION, + reason="Installed google-api-core lacks auto-instrumentation suppression", +) +def test_auto_instrumentation_suppression_sync(span_exporter): + """Verifies that upstream gRPC client auto-instrumentation spans are suppressed in sync calls. + + When users enable global GrpcInstrumentorClient().instrument(), upstream injects a + generic interceptor into channel creation. This test verifies that our downstream + suppression interceptor prevents the duplicate, bare-bones upstream span from being + emitted while preserving the full Google Cloud SDK T4 span. + """ + grpc_instrumentation = pytest.importorskip("opentelemetry.instrumentation.grpc") + + exporter, provider = span_exporter + instrumentor = grpc_instrumentation.GrpcInstrumentorClient() + instrumentor.instrument(tracer_provider=provider) + try: + options = ClientOptions( + api_endpoint="localhost:7469", + tracer_provider=provider, + ) + with mock.patch.dict( + os.environ, {"GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED": "true"} + ): + channel = grpc.insecure_channel("localhost:7469") + transport = EchoClient.get_transport_class("grpc")( + channel=channel, + client_options=options, + credentials=ga_credentials.AnonymousCredentials(), + ) + client = EchoClient(transport=transport, client_options=options) + response = client.echo( + showcase.EchoRequest(content="suppress sync duplicate") + ) + assert response.content == "suppress sync duplicate" + + spans = exporter.get_finished_spans() + assert len(spans) == 2 + + root_spans = [s for s in spans if s.parent is None] + child_spans = [s for s in spans if s.parent is not None] + assert len(root_spans) == 1 + assert len(child_spans) == 1 + + root_span = root_spans[0] + child_span = child_spans[0] + + assert root_span.name == "google.showcase.v1beta1.Echo/Echo" + assert child_span.name == "google.showcase.v1beta1.Echo/Echo" + assert child_span.parent.span_id == root_span.context.span_id + + assert child_span.attributes.get("rpc.system.name") == "grpc" + assert child_span.attributes.get("server.address") == "localhost" + assert child_span.attributes.get("server.port") == 7469 + assert child_span.attributes.get("url.domain") == "googleapis.com" + assert child_span.attributes.get("rpc.response.status_code") == "OK" + finally: + instrumentor.uninstrument() + + +@pytest.mark.asyncio +@pytest.mark.skipif( + not HAS_AUTO_INSTRUMENTATION_SUPPRESSION, + reason="Installed google-api-core lacks auto-instrumentation suppression", +) +async def test_auto_instrumentation_suppression_async(span_exporter): + """Verifies that upstream gRPC client auto-instrumentation spans are suppressed in async calls. + + When users enable global GrpcAioInstrumentorClient().instrument(), upstream injects a + generic interceptor into aio channel creation. This test verifies that our downstream + suppression interceptor prevents the duplicate, bare-bones upstream span from being + emitted while preserving the full Google Cloud SDK T4 span. + """ + grpc_instrumentation = pytest.importorskip("opentelemetry.instrumentation.grpc") + + exporter, provider = span_exporter + instrumentor = grpc_instrumentation.GrpcAioInstrumentorClient() + instrumentor.instrument(tracer_provider=provider) + try: + options = ClientOptions( + api_endpoint="localhost:7469", + tracer_provider=provider, + ) + with mock.patch.dict( + os.environ, {"GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED": "true"} + ): + channel = grpc.aio.insecure_channel("localhost:7469") + transport = EchoAsyncClient.get_transport_class("grpc_asyncio")( + channel=channel, + client_options=options, + credentials=ga_credentials.AnonymousCredentials(), + ) + client = EchoAsyncClient(transport=transport, client_options=options) + response = await client.echo( + showcase.EchoRequest(content="suppress async duplicate") + ) + assert response.content == "suppress async duplicate" + + spans = exporter.get_finished_spans() + assert len(spans) == 2 + + root_spans = [s for s in spans if s.parent is None] + child_spans = [s for s in spans if s.parent is not None] + assert len(root_spans) == 1 + assert len(child_spans) == 1 + + root_span = root_spans[0] + child_span = child_spans[0] + + assert root_span.name == "google.showcase.v1beta1.Echo/Echo" + assert child_span.name == "google.showcase.v1beta1.Echo/Echo" + assert child_span.parent.span_id == root_span.context.span_id + + assert child_span.attributes.get("rpc.system.name") == "grpc" + assert child_span.attributes.get("server.address") == "localhost" + assert child_span.attributes.get("server.port") == 7469 + assert child_span.attributes.get("url.domain") == "googleapis.com" + assert child_span.attributes.get("rpc.response.status_code") == "OK" + finally: + instrumentor.uninstrument() diff --git a/packages/google-api-core/google/api_core/_observability.py b/packages/google-api-core/google/api_core/_observability.py index cddc579b99bf..89e33b25f28c 100644 --- a/packages/google-api-core/google/api_core/_observability.py +++ b/packages/google-api-core/google/api_core/_observability.py @@ -18,12 +18,20 @@ from __future__ import annotations +import contextlib import urllib.parse from typing import TYPE_CHECKING, Any, Callable, Sequence from google.api_core import _feature_gating_helpers from google.api_core.client_options import ClientOptions +try: + from opentelemetry.instrumentation.utils import ( + suppress_instrumentation as _suppress_instrumentation, + ) +except ImportError: + _suppress_instrumentation = contextlib.nullcontext + if TYPE_CHECKING: # flake8: grpc, trace, and ClientInterceptor are imported only for static analysis and type annotations # The 'noqa: F401' comment avoids flake8 "imported but not used" errors. @@ -223,6 +231,148 @@ def _get_tracer_provider( return None +try: + # flake8: `grpc` imported conditionally for runtime interceptor base classes + import grpc # noqa: F811 + from grpc import ( + StreamStreamClientInterceptor as _SyncStreamStreamClientInterceptor, + ) + from grpc import ( + StreamUnaryClientInterceptor as _SyncStreamUnaryClientInterceptor, + ) + from grpc import ( + UnaryStreamClientInterceptor as _SyncUnaryStreamClientInterceptor, + ) + from grpc import ( + UnaryUnaryClientInterceptor as _SyncUnaryUnaryClientInterceptor, + ) +except ImportError: # pragma: NO COVER + # mypy: Fallback sentinel when optional grpc is not installed in the environment. + grpc = None # type: ignore[assignment] + + # mypy: Fallback dummy classes when optional grpc is not installed. + # Four distinct empty classes avoid duplicate base class 'object' TypeError at runtime. + class _SyncUnaryUnaryClientInterceptor: # type: ignore[no-redef] + pass + + class _SyncUnaryStreamClientInterceptor: # type: ignore[no-redef] + pass + + class _SyncStreamUnaryClientInterceptor: # type: ignore[no-redef] + pass + + class _SyncStreamStreamClientInterceptor: # type: ignore[no-redef] + pass + + +try: + from grpc.aio import ( + StreamStreamClientInterceptor as _AsyncStreamStreamClientInterceptor, + ) + from grpc.aio import ( + StreamUnaryClientInterceptor as _AsyncStreamUnaryClientInterceptor, + ) + from grpc.aio import ( + UnaryStreamClientInterceptor as _AsyncUnaryStreamClientInterceptor, + ) + from grpc.aio import ( + UnaryUnaryClientInterceptor as _AsyncUnaryUnaryClientInterceptor, + ) +except ImportError: # pragma: NO COVER + # mypy: Fallback dummy classes when optional grpc is not installed. + # Four distinct empty classes avoid duplicate base class 'object' TypeError at runtime. + class _AsyncUnaryUnaryClientInterceptor: # type: ignore[no-redef] + pass + + class _AsyncUnaryStreamClientInterceptor: # type: ignore[no-redef] + pass + + class _AsyncStreamUnaryClientInterceptor: # type: ignore[no-redef] + pass + + class _AsyncStreamStreamClientInterceptor: # type: ignore[no-redef] + pass + + +class _SuppressingClientInterceptor( + _SyncUnaryUnaryClientInterceptor, + _SyncUnaryStreamClientInterceptor, + _SyncStreamUnaryClientInterceptor, + _SyncStreamStreamClientInterceptor, +): + """Client interceptor that suppresses redundant downstream generic auto-instrumentation spans. + + When users enable global OpenTelemetry gRPC auto-instrumentation (e.g. GrpcInstrumentorClient), + the underlying channel is wrapped with generic interceptors that emit bare-bones spans. + This interceptor wraps downstream calls in OpenTelemetry's native `suppress_instrumentation` + context manager so that our feature-rich Google Cloud SDK T4 span is emitted while + redundant generic downstream spans are bypassed. + """ + + def __init__(self) -> None: + self._is_otel_interceptor = True + + def intercept_unary_unary( + self, continuation: Any, client_call_details: Any, request: Any + ) -> Any: + with _suppress_instrumentation(): + return continuation(client_call_details, request) + + def intercept_unary_stream( + self, continuation: Any, client_call_details: Any, request: Any + ) -> Any: + with _suppress_instrumentation(): + return continuation(client_call_details, request) + + def intercept_stream_unary( + self, continuation: Any, client_call_details: Any, request_iterator: Any + ) -> Any: + with _suppress_instrumentation(): + return continuation(client_call_details, request_iterator) + + def intercept_stream_stream( + self, continuation: Any, client_call_details: Any, request_iterator: Any + ) -> Any: + with _suppress_instrumentation(): + return continuation(client_call_details, request_iterator) + + +class _AsyncSuppressingClientInterceptor( + _AsyncUnaryUnaryClientInterceptor, + _AsyncUnaryStreamClientInterceptor, + _AsyncStreamUnaryClientInterceptor, + _AsyncStreamStreamClientInterceptor, +): + """AsyncIO client interceptor that suppresses redundant downstream generic auto-instrumentation spans.""" + + def __init__(self) -> None: + self._is_otel_interceptor = True + + async def intercept_unary_unary( + self, continuation: Any, client_call_details: Any, request: Any + ) -> Any: + with _suppress_instrumentation(): + return await continuation(client_call_details, request) + + async def intercept_unary_stream( + self, continuation: Any, client_call_details: Any, request: Any + ) -> Any: + with _suppress_instrumentation(): + return await continuation(client_call_details, request) + + async def intercept_stream_unary( + self, continuation: Any, client_call_details: Any, request_iterator: Any + ) -> Any: + with _suppress_instrumentation(): + return await continuation(client_call_details, request_iterator) + + async def intercept_stream_stream( + self, continuation: Any, client_call_details: Any, request_iterator: Any + ) -> Any: + with _suppress_instrumentation(): + return await continuation(client_call_details, request_iterator) + + def get_otel_interceptor( client_options: ClientOptions | dict[str, Any] | None = None, ) -> Callable[[grpc.Channel], grpc.Channel] | None: @@ -249,8 +399,15 @@ def get_otel_interceptor( request_hook=request_hook, response_hook=_grpc_client_response_hook, ) + suppressor = _SuppressingClientInterceptor() def otel_interceptor(channel: grpc.Channel) -> grpc.Channel: + if grpc is not None: + try: + channel = grpc.intercept_channel(channel, suppressor) + except (AttributeError, TypeError): # pragma: NO COVER + # Fallback if grpc lacks intercept_channel or custom channel type rejects it. + pass return otel_grpc.intercept_channel(channel, interceptor) otel_interceptor._is_otel_interceptor = True # type: ignore[attr-defined] @@ -279,11 +436,21 @@ def get_otel_async_interceptor( endpoint_attrs = _extract_endpoint_attributes(client_options) request_hook = _make_grpc_client_request_hook(endpoint_attrs) - return otel_grpc.aio_client_interceptors( + raw_interceptors = otel_grpc.aio_client_interceptors( tracer_provider=_get_tracer_provider(client_options), request_hook=request_hook, response_hook=_grpc_client_response_hook, ) + interceptors = ( + list(raw_interceptors) + if isinstance(raw_interceptors, (list, tuple)) + else [raw_interceptors] + ) + interceptors.append(_AsyncSuppressingClientInterceptor()) + for interceptor in interceptors: + interceptor._is_otel_interceptor = True # type: ignore[attr-defined] + + return interceptors _TRACE_CONTEXT_PROPAGATOR: Any = None 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 f0cdd1905a9a..42834453a1d3 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 @@ -338,25 +338,36 @@ def apply_channel_interceptors( ("intercept_stream_unary", "_stream_unary_interceptors"), ("intercept_stream_stream", "_stream_stream_interceptors"), ) + # Track insertion index per interceptor list for OpenTelemetry interceptors. + # OpenTelemetry interceptors and downstream suppressors are inserted at the + # beginning of the channel's interceptor pipeline in order, so that our + # Google SDK spans wrap outside and downstream generic auto-instrumentation + # (e.g. GrpcAioInstrumentorClient) is suppressed. Other interceptors are appended. + otel_insert_indices: dict[str, int] = {} for interceptor in interceptors: - matched = False - for method_name, attr_name in mapping: - if hasattr(interceptor, method_name) and hasattr(channel, attr_name): - target_list = getattr(channel, attr_name) - if isinstance(target_list, list): - if interceptor not in target_list: - target_list.append(interceptor) - matched = True - elif hasattr(target_list, "append"): + is_otel = getattr(interceptor, "_is_otel_interceptor", False) is True + + targets = [ + attr_name + for method_name, attr_name in mapping + if hasattr(interceptor, method_name) and hasattr(channel, attr_name) + ] + if not targets and hasattr(channel, "_unary_unary_interceptors"): + targets = ["_unary_unary_interceptors"] + + for attr_name in targets: + target_list = getattr(channel, attr_name) + if isinstance(target_list, list): + if interceptor in target_list: + continue + if is_otel: + idx = otel_insert_indices.get(attr_name, 0) + target_list.insert(idx, interceptor) + otel_insert_indices[attr_name] = idx + 1 + else: target_list.append(interceptor) - matched = True - if not matched and hasattr(channel, "_unary_unary_interceptors"): - unary_interceptors = channel._unary_unary_interceptors - if isinstance(unary_interceptors, list): - if interceptor not in unary_interceptors: - unary_interceptors.append(interceptor) - elif hasattr(unary_interceptors, "append"): - unary_interceptors.append(interceptor) + elif hasattr(target_list, "append"): + target_list.append(interceptor) return channel 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 915de50297b3..a0b458c735e2 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 @@ -832,3 +832,64 @@ class CustomInterceptor: result = grpc_helpers_async.apply_channel_interceptors(channel, [interceptor]) assert result is channel mock_append.append.assert_called_once_with(interceptor) + + +def test_apply_channel_interceptors_otel_prepended_in_order(): + """Proves that interceptors with _is_otel_interceptor=True are prepended in order.""" + existing_interceptor = mock.Mock() + otel_interceptor1 = mock.Mock(spec=["intercept_unary_unary"]) + otel_interceptor1._is_otel_interceptor = True + otel_interceptor2 = mock.Mock(spec=["intercept_unary_unary"]) + otel_interceptor2._is_otel_interceptor = True + regular_interceptor = mock.Mock(spec=["intercept_unary_unary"]) + + channel = mock.Mock() + channel._unary_unary_interceptors = [existing_interceptor] + + result = grpc_helpers_async.apply_channel_interceptors( + channel, [regular_interceptor, otel_interceptor1, otel_interceptor2] + ) + assert result is channel + # OTel interceptors should be inserted at index 0 in order, regular interceptor appended + assert channel._unary_unary_interceptors == [ + otel_interceptor1, + otel_interceptor2, + existing_interceptor, + regular_interceptor, + ] + + +def test_apply_channel_interceptors_fallback_otel_prepended_in_order(): + """Proves that fallback interceptors with _is_otel_interceptor=True are prepended in order. + + Standard gRPC interceptors implement canonical methods such as + `intercept_unary_unary` or `intercept_stream_stream`. When an interceptor + does not define any of the 4 standard methods, `apply_channel_interceptors` + falls back to attaching it directly to `channel._unary_unary_interceptors`. + + This test verifies that if such a non-standard fallback interceptor is flagged + with `_is_otel_interceptor = True`, it is prepended at index 0 in pipeline order + rather than appended to the end. + """ + + class FallbackOtelInterceptor: + _is_otel_interceptor = True + + existing_interceptor = mock.Mock() + otel1 = FallbackOtelInterceptor() + otel2 = FallbackOtelInterceptor() + regular = mock.Mock(spec=[]) + + channel = mock.Mock(spec=["_unary_unary_interceptors"]) + channel._unary_unary_interceptors = [existing_interceptor] + + result = grpc_helpers_async.apply_channel_interceptors( + channel, [regular, otel1, otel2] + ) + assert result is channel + assert channel._unary_unary_interceptors == [ + otel1, + otel2, + existing_interceptor, + regular, + ] diff --git a/packages/google-api-core/tests/conftest.py b/packages/google-api-core/tests/conftest.py index 664ef7f95895..1fa774ded176 100644 --- a/packages/google-api-core/tests/conftest.py +++ b/packages/google-api-core/tests/conftest.py @@ -64,3 +64,22 @@ def mock_otel(monkeypatch): tracer=mock_tracer, span=mock_span, ) + + +@pytest.fixture +def mock_otel_grpc(monkeypatch): + """Provides a mocked opentelemetry.instrumentation.grpc module environment. + + Unlike `mock_otel` (which mocks `opentelemetry.trace` to test span generation + in client methods), this fixture simulates the installation of the gRPC + auto-instrumentation package (`opentelemetry.instrumentation.grpc`) in `sys.modules`. + This allows testing gRPC channel interceptor resolution and capability checks. + """ + mock_otel = mock.Mock() + mock_grpc = mock_otel.instrumentation.grpc + monkeypatch.setitem(sys.modules, "opentelemetry", mock_otel) + monkeypatch.setitem( + sys.modules, "opentelemetry.instrumentation", mock_otel.instrumentation + ) + monkeypatch.setitem(sys.modules, "opentelemetry.instrumentation.grpc", mock_grpc) + return mock_grpc diff --git a/packages/google-api-core/tests/unit/test_observability.py b/packages/google-api-core/tests/unit/test_observability.py index 90ac4acb3ba3..35e09561af95 100644 --- a/packages/google-api-core/tests/unit/test_observability.py +++ b/packages/google-api-core/tests/unit/test_observability.py @@ -40,23 +40,11 @@ def test_is_otel_capabilities_enabled_otel_missing(monkeypatch): assert not _observability.is_otel_capabilities_enabled() -def test_is_otel_capabilities_enabled_otel_installed(monkeypatch): +def test_is_otel_capabilities_enabled_otel_installed(monkeypatch, mock_otel_grpc): """Proves that is_otel_capabilities_enabled returns True when tracing is enabled and OpenTelemetry gRPC instrumentation is installed. """ monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") - - mock_otel = mock.Mock() - mock_otel_grpc = mock_otel.instrumentation.grpc - - monkeypatch.setitem(sys.modules, "opentelemetry", mock_otel) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation", mock_otel.instrumentation - ) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation.grpc", mock_otel_grpc - ) - assert _observability.is_otel_capabilities_enabled() @@ -74,23 +62,13 @@ def test_is_otel_capabilities_enabled_experimental_requires_env_var(monkeypatch) _observability.is_otel_capabilities_enabled(options) -def test_is_otel_capabilities_enabled_experimental_enabled_with_config(monkeypatch): +def test_is_otel_capabilities_enabled_experimental_enabled_with_config( + monkeypatch, mock_otel_grpc +): """Proves that when GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED=true and tracer_provider is supplied via client_options, is_otel_capabilities_enabled returns True. """ monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") - - mock_otel = mock.Mock() - mock_otel_grpc = mock_otel.instrumentation.grpc - - monkeypatch.setitem(sys.modules, "opentelemetry", mock_otel) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation", mock_otel.instrumentation - ) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation.grpc", mock_otel_grpc - ) - options = ClientOptions(tracer_provider=mock.Mock()) assert _observability.is_otel_capabilities_enabled(options) @@ -135,22 +113,38 @@ def test_get_tracer_provider_dict_config(): assert _observability._get_tracer_provider(options) is mock_tracer_provider -def test_get_otel_interceptor_disabled(monkeypatch): - """Proves that get_otel_interceptor returns None when tracing is disabled.""" +@pytest.mark.parametrize( + "interceptor_getter", + [ + _observability.get_otel_interceptor, + _observability.get_otel_async_interceptor, + ], + ids=["sync_interceptor", "async_interceptor"], +) +def test_get_otel_interceptor_disabled(monkeypatch, interceptor_getter): + """Proves that interceptor getters return None when tracing is disabled.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "false") - assert _observability.get_otel_interceptor() is None + assert interceptor_getter() is None -def test_get_otel_interceptor_otel_missing(monkeypatch): - """Proves that get_otel_interceptor returns None when OpenTelemetry gRPC +@pytest.mark.parametrize( + "interceptor_getter", + [ + _observability.get_otel_interceptor, + _observability.get_otel_async_interceptor, + ], + ids=["sync_interceptor", "async_interceptor"], +) +def test_get_otel_interceptor_otel_missing(monkeypatch, interceptor_getter): + """Proves that interceptor getters return None when OpenTelemetry gRPC instrumentation is not installed. """ monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") monkeypatch.setitem(sys.modules, "opentelemetry.instrumentation.grpc", None) - assert _observability.get_otel_interceptor() is None + assert interceptor_getter() is None -def test_get_otel_interceptor_enabled(monkeypatch): +def test_get_otel_interceptor_enabled(monkeypatch, mock_otel_grpc): """Proves that get_otel_interceptor creates a synchronous OpenTelemetry client interceptor with the resolved tracer provider and returns a channel-intercepting callable. """ @@ -160,22 +154,11 @@ def test_get_otel_interceptor_enabled(monkeypatch): mock_raw_channel = mock.Mock(name="raw_channel") mock_wrapped_channel = mock.Mock(name="wrapped_channel") - - mock_otel = mock.Mock() - mock_otel_grpc = mock_otel.instrumentation.grpc mock_interceptor = mock.Mock(name="otel_interceptor") mock_otel_grpc.client_interceptor.return_value = mock_interceptor mock_otel_grpc.intercept_channel.return_value = mock_wrapped_channel - monkeypatch.setitem(sys.modules, "opentelemetry", mock_otel) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation", mock_otel.instrumentation - ) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation.grpc", mock_otel_grpc - ) - interceptor = _observability.get_otel_interceptor(client_options=options) assert callable(interceptor) @@ -192,12 +175,14 @@ def test_get_otel_interceptor_enabled(monkeypatch): result = interceptor(mock_raw_channel) assert result is mock_wrapped_channel - mock_otel_grpc.intercept_channel.assert_called_once_with( - mock_raw_channel, mock_interceptor - ) + assert mock_otel_grpc.intercept_channel.call_count == 1 + chan_arg, interceptor_arg = mock_otel_grpc.intercept_channel.call_args[0] + assert interceptor_arg is mock_interceptor -def test_get_otel_interceptor_with_apply_channel_interceptors(monkeypatch): +def test_get_otel_interceptor_with_apply_channel_interceptors( + monkeypatch, mock_otel_grpc +): """Proves that get_otel_interceptor integrates seamlessly into apply_channel_interceptors.""" pytest.importorskip("grpc") from google.api_core import grpc_helpers @@ -208,22 +193,11 @@ def test_get_otel_interceptor_with_apply_channel_interceptors(monkeypatch): mock_raw_channel = mock.Mock(name="raw_channel") mock_wrapped_channel = mock.Mock(name="wrapped_channel") - - mock_otel = mock.Mock() - mock_otel_grpc = mock_otel.instrumentation.grpc mock_interceptor = mock.Mock(name="otel_interceptor") mock_otel_grpc.client_interceptor.return_value = mock_interceptor mock_otel_grpc.intercept_channel.return_value = mock_wrapped_channel - monkeypatch.setitem(sys.modules, "opentelemetry", mock_otel) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation", mock_otel.instrumentation - ) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation.grpc", mock_otel_grpc - ) - otel_interceptor = _observability.get_otel_interceptor(client_options=options) assert callable(otel_interceptor) @@ -231,27 +205,36 @@ def test_get_otel_interceptor_with_apply_channel_interceptors(monkeypatch): mock_raw_channel, interceptors=[otel_interceptor] ) assert result is mock_wrapped_channel - mock_otel_grpc.intercept_channel.assert_called_once_with( - mock_raw_channel, mock_interceptor - ) - - -def test_get_otel_async_interceptor_disabled(monkeypatch): - """Proves that get_otel_async_interceptor returns None when tracing is disabled.""" - monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "false") - assert _observability.get_otel_async_interceptor() is None + assert mock_otel_grpc.intercept_channel.call_count == 1 + chan_arg, interceptor_arg = mock_otel_grpc.intercept_channel.call_args[0] + assert interceptor_arg is mock_interceptor -def test_get_otel_async_interceptor_otel_missing(monkeypatch): - """Proves that get_otel_async_interceptor returns None when OpenTelemetry gRPC - instrumentation is not installed. - """ +@pytest.mark.parametrize("error_cls", [AttributeError, TypeError]) +def test_get_otel_interceptor_grpc_intercept_channel_fallback( + monkeypatch, mock_otel_grpc, error_cls +): + """Proves that otel_interceptor falls back gracefully if grpc.intercept_channel raises AttributeError or TypeError.""" + if _observability.grpc is None: + pytest.skip("grpc is not installed") monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") - monkeypatch.setitem(sys.modules, "opentelemetry.instrumentation.grpc", None) - assert _observability.get_otel_async_interceptor() is None + options = ClientOptions() + mock_raw_channel = mock.Mock(name="raw_channel") + mock_wrapped_channel = mock.Mock(name="wrapped_channel") + mock_otel_grpc.client_interceptor.return_value = mock.Mock() + mock_otel_grpc.intercept_channel.return_value = mock_wrapped_channel + otel_interceptor = _observability.get_otel_interceptor(client_options=options) + with mock.patch.object( + _observability.grpc, + "intercept_channel", + side_effect=error_cls("Mock intercept error"), + ): + result = otel_interceptor(mock_raw_channel) + assert result is mock_wrapped_channel -def test_get_otel_async_interceptor_enabled(monkeypatch): + +def test_get_otel_async_interceptor_enabled(monkeypatch, mock_otel_grpc): """Proves that get_otel_async_interceptor instantiates and returns asynchronous OpenTelemetry client interceptors with the resolved tracer provider. """ @@ -260,21 +243,13 @@ def test_get_otel_async_interceptor_enabled(monkeypatch): options = ClientOptions(tracer_provider=mock_tracer_provider) mock_async_interceptors = [mock.Mock(name="otel_async_interceptor")] - - mock_otel = mock.Mock() - mock_otel_grpc = mock_otel.instrumentation.grpc mock_otel_grpc.aio_client_interceptors.return_value = mock_async_interceptors - monkeypatch.setitem(sys.modules, "opentelemetry", mock_otel) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation", mock_otel.instrumentation - ) - monkeypatch.setitem( - sys.modules, "opentelemetry.instrumentation.grpc", mock_otel_grpc - ) - result = _observability.get_otel_async_interceptor(client_options=options) - assert result is mock_async_interceptors + assert mock_async_interceptors[0] in result + assert any( + isinstance(i, _observability._AsyncSuppressingClientInterceptor) for i in result + ) mock_otel_grpc.aio_client_interceptors.assert_called_once_with( tracer_provider=mock_tracer_provider, request_hook=mock.ANY, @@ -288,6 +263,36 @@ def test_get_otel_async_interceptor_enabled(monkeypatch): mock_span.set_attribute.assert_any_call("url.domain", "googleapis.com") +def test_get_otel_async_interceptor_single_return(monkeypatch, mock_otel_grpc): + """Verifies that a single interceptor returned by upstream OTel is safely wrapped into a list.""" + # Step 1: Opt-in to experimental SDK tracing via environment variable. + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + + # Step 2: Configure the mock OpenTelemetry gRPC instrumentation fixture. + # While upstream normally returns a sequence (list or tuple), third-party + # instrumentation or older versions may return a lone interceptor object. + # We simulate this edge case by returning a single mock interceptor. + single_interceptor = mock.Mock(name="single_otel_interceptor") + mock_otel_grpc.aio_client_interceptors.return_value = single_interceptor + + # Step 3: Call get_otel_async_interceptor with standard ClientOptions. + options = ClientOptions() + result = _observability.get_otel_async_interceptor(client_options=options) + + # Step 4: Validate list normalization and downstream suppression interceptor injection. + # The result must be a 2-item list: + # [0]: The single upstream interceptor wrapped in a list. + # [1]: Our internal _AsyncSuppressingClientInterceptor appended to prevent redundant wire spans. + assert isinstance(result, list) + assert len(result) == 2 + assert result[0] is single_interceptor + assert isinstance(result[1], _observability._AsyncSuppressingClientInterceptor) + + # Step 5: Verify the internal flag marking that channel transports use to avoid duplicate wrapping. + assert getattr(result[0], "_is_otel_interceptor", False) is True + assert getattr(result[1], "_is_otel_interceptor", False) is True + + @pytest.mark.parametrize( "client_options,expected_attrs", [ @@ -1302,3 +1307,78 @@ def test_trace_http_request_initialization_fails_open(monkeypatch): assert isinstance(ctx, _observability._TraceContext) assert ctx._span is None ctx.record_response(mock.Mock()) + + +@pytest.mark.parametrize( + "method_name", + [ + "intercept_unary_unary", + "intercept_unary_stream", + "intercept_stream_unary", + "intercept_stream_stream", + ], + ids=["unary_unary", "unary_stream", "stream_unary", "stream_stream"], +) +def test_sync_suppressing_interceptor_invocations(method_name): + """Verifies that _SuppressingClientInterceptor wraps downstream calls in suppress_instrumentation. + + When global OpenTelemetry gRPC auto-instrumentation is enabled, it emits + generic, un-enriched attempt spans. We wrap downstream calls in + suppress_instrumentation so only our enriched Google Cloud SDK T4 span is emitted. + """ + suppressor = _observability._SuppressingClientInterceptor() + assert getattr(suppressor, "_is_otel_interceptor", False) is True + + with mock.patch( + "google.api_core._observability._suppress_instrumentation" + ) as mock_suppress: + mock_continuation = mock.Mock(return_value="expected_response") + interceptor_method = getattr(suppressor, method_name) + result = interceptor_method( + mock_continuation, "call_details", "request_or_iterator" + ) + + assert result == "expected_response" + mock_continuation.assert_called_once_with("call_details", "request_or_iterator") + mock_suppress.assert_called_once() + + +@pytest.mark.parametrize( + "method_name", + [ + "intercept_unary_unary", + "intercept_unary_stream", + "intercept_stream_unary", + "intercept_stream_stream", + ], + ids=["unary_unary", "unary_stream", "stream_unary", "stream_stream"], +) +def test_async_suppressing_interceptor_invocations(method_name): + """Verifies that _AsyncSuppressingClientInterceptor wraps downstream async calls in suppress_instrumentation. + + When global OpenTelemetry gRPC auto-instrumentation is enabled, it emits + generic, un-enriched attempt spans. We wrap downstream async calls in + suppress_instrumentation so only our enriched Google Cloud SDK T4 span is emitted. + """ + import asyncio + + suppressor = _observability._AsyncSuppressingClientInterceptor() + assert getattr(suppressor, "_is_otel_interceptor", False) is True + + async def _run_test(): + with mock.patch( + "google.api_core._observability._suppress_instrumentation" + ) as mock_suppress: + mock_continuation = mock.AsyncMock(return_value="expected_response") + interceptor_method = getattr(suppressor, method_name) + result = await interceptor_method( + mock_continuation, "call_details", "request_or_iterator" + ) + + assert result == "expected_response" + mock_continuation.assert_called_once_with( + "call_details", "request_or_iterator" + ) + mock_suppress.assert_called_once() + + asyncio.run(_run_test())