From 95ca42a7eead9b57ef07a35dcd42b4f675d1801a Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 3 Sep 2026 09:50:23 -0400 Subject: [PATCH 01/11] feat(gapic): add OpenTelemetry T3 client method span wrapping --- .../google/api_core/gapic_v1/method.py | 39 ++++++- .../tests/unit/gapic/test_method.py | 100 ++++++++++++++++++ 2 files changed, 138 insertions(+), 1 deletion(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index ecd54d0aef62..aceecdd5e96c 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -22,7 +22,7 @@ import functools from typing import List, Tuple -from google.api_core import grpc_helpers +from google.api_core import _observability, grpc_helpers from google.api_core.gapic_v1 import client_info from google.api_core.timeout import TimeToDeadlineTimeout @@ -125,6 +125,8 @@ class _GapicCallable(object): additional metadata will be passed to the RPC method. """ + _is_tracing_supported = True + def __init__( self, target, @@ -186,6 +188,41 @@ def __call__( if self._compression is not None: kwargs["compression"] = compression + if _observability.is_otel_capabilities_enabled(): + try: + from opentelemetry import trace + + tracer = trace.get_tracer("google.api_core") + raw_method = getattr(self._target, "_method", None) + if raw_method and isinstance(raw_method, (str, bytes)): + if isinstance(raw_method, bytes): + raw_method = raw_method.decode("utf-8") + method_str = raw_method.lstrip("/") + service, _, method = method_str.rpartition("/") + span_name = method_str + else: + service = "google.api_core" + method = getattr(self._target, "__name__", "call") + span_name = f"{service}/{method}" + + with tracer.start_as_current_span( + span_name, + kind=trace.SpanKind.CLIENT, + attributes={ + "rpc.system": "grpc", + "rpc.service": service, + "rpc.method": method, + }, + ) as span: + try: + return wrapped_func(*args, **kwargs) + except Exception as exc: + span.record_exception(exc) + span.set_status(trace.StatusCode.ERROR, str(exc)) + raise + except ImportError: + pass + return wrapped_func(*args, **kwargs) diff --git a/packages/google-api-core/tests/unit/gapic/test_method.py b/packages/google-api-core/tests/unit/gapic/test_method.py index fbe7f2a5f0f1..6f5ab8c1bd9a 100644 --- a/packages/google-api-core/tests/unit/gapic/test_method.py +++ b/packages/google-api-core/tests/unit/gapic/test_method.py @@ -13,6 +13,7 @@ # limitations under the License. import datetime +import sys from unittest import mock import pytest @@ -346,3 +347,102 @@ def test_wrap_method_with_call_not_supported(): def test__deduplicate_metadata_tokens(headers, expected): dedup = google.api_core.gapic_v1.method._deduplicate_metadata_tokens assert dedup(*headers) == expected + + +def test_wrap_method_otel_tracing_disabled(monkeypatch): + """Proves that when OpenTelemetry tracing is disabled, no span is created.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "false") + mock_target = mock.Mock(return_value="success") + wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) + + with mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=False, + ): + assert wrapped() == "success" + mock_target.assert_called_once() + + +def test_wrap_method_otel_tracing_enabled_success(monkeypatch): + """Proves that when OpenTelemetry tracing is enabled, a T3 client span is started.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="success") + mock_target._method = ( + "/google.cloud.secretmanager.v1.SecretManagerService/ListSecrets" + ) + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_trace = mock.Mock() + mock_trace.get_tracer.return_value = mock_tracer + mock_trace.SpanKind.CLIENT = "CLIENT" + + wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + result = wrapped() + + assert result == "success" + mock_tracer.start_as_current_span.assert_called_once_with( + "google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.cloud.secretmanager.v1.SecretManagerService", + "rpc.method": "ListSecrets", + }, + ) + + +def test_wrap_method_otel_tracing_enabled_error(monkeypatch): + """Proves that when an RPC fails, the T3 client span records the exception and error status.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + err = RuntimeError("gRPC connection reset") + mock_target = mock.Mock(side_effect=err) + mock_target._method = ( + "/google.cloud.secretmanager.v1.SecretManagerService/ListSecrets" + ) + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_trace = mock.Mock() + mock_trace.get_tracer.return_value = mock_tracer + mock_trace.SpanKind.CLIENT = "CLIENT" + mock_trace.StatusCode.ERROR = "ERROR" + + wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + with pytest.raises(RuntimeError): + wrapped() + + mock_span.record_exception.assert_called_once_with(err) + mock_span.set_status.assert_called_once_with("ERROR", str(err)) From 4cd28f0bb9b24c16ea7a0a01a49a12d2ab85029b Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 3 Sep 2026 10:43:15 -0400 Subject: [PATCH 02/11] chore(gapic): remove _is_tracing_supported dummy class variable --- packages/google-api-core/google/api_core/gapic_v1/method.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index aceecdd5e96c..d2022756f3b2 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -125,8 +125,6 @@ class _GapicCallable(object): additional metadata will be passed to the RPC method. """ - _is_tracing_supported = True - def __init__( self, target, From 9b46ef76c40c7b187a0ea36ce93149afad631b01 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Fri, 4 Sep 2026 11:18:49 -0400 Subject: [PATCH 03/11] test(gapic): achieve 100% branch and statement coverage for OTel T3 method tracing --- .../google/api_core/gapic_v1/method.py | 3 +- .../tests/unit/gapic/test_method.py | 212 ++++++++++++++++++ 2 files changed, 214 insertions(+), 1 deletion(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index d2022756f3b2..dceb094f5568 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -218,7 +218,8 @@ def __call__( span.record_exception(exc) span.set_status(trace.StatusCode.ERROR, str(exc)) raise - except ImportError: + # If OpenTelemetry cannot be imported in the current environment, continue without tracing. + except ImportError: # pragma: NO COVER pass return wrapped_func(*args, **kwargs) diff --git a/packages/google-api-core/tests/unit/gapic/test_method.py b/packages/google-api-core/tests/unit/gapic/test_method.py index 6f5ab8c1bd9a..c337a31fe6f6 100644 --- a/packages/google-api-core/tests/unit/gapic/test_method.py +++ b/packages/google-api-core/tests/unit/gapic/test_method.py @@ -446,3 +446,215 @@ def test_wrap_method_otel_tracing_enabled_error(monkeypatch): mock_span.record_exception.assert_called_once_with(err) mock_span.set_status.assert_called_once_with("ERROR", str(err)) + + +def test_wrap_method_otel_tracing_bytes_method(monkeypatch): + """Proves that when raw _method is bytes, it is decoded properly to utf-8.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="success") + mock_target._method = ( + b"/google.cloud.secretmanager.v1.SecretManagerService/ListSecrets" + ) + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_trace = mock.Mock() + mock_trace.get_tracer.return_value = mock_tracer + mock_trace.SpanKind.CLIENT = "CLIENT" + + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, default_timeout=60 + ) + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + result = wrapped() + + assert result == "success" + mock_tracer.start_as_current_span.assert_called_once_with( + "google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.cloud.secretmanager.v1.SecretManagerService", + "rpc.method": "ListSecrets", + }, + ) + + +def test_wrap_method_otel_tracing_fallback_with_name(monkeypatch): + """Proves that when raw _method is absent, fallback uses target.__name__.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + + def custom_rpc(*args, **kwargs): + return "success" + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_trace = mock.Mock() + mock_trace.get_tracer.return_value = mock_tracer + mock_trace.SpanKind.CLIENT = "CLIENT" + + wrapped = google.api_core.gapic_v1.method.wrap_method(custom_rpc) + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + result = wrapped() + + assert result == "success" + mock_tracer.start_as_current_span.assert_called_once_with( + "google.api_core/custom_rpc", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.api_core", + "rpc.method": "custom_rpc", + }, + ) + + +def test_wrap_method_otel_tracing_fallback_without_name(monkeypatch): + """Proves that when raw _method is absent and target has no explicit __name__, + fallback uses target class name assigned by error wrapper. + """ + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + + class TargetWithoutName: + def __call__(self, *args, **kwargs): + return "success" + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_trace = mock.Mock() + mock_trace.get_tracer.return_value = mock_tracer + mock_trace.SpanKind.CLIENT = "CLIENT" + + wrapped = google.api_core.gapic_v1.method.wrap_method(TargetWithoutName()) + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + result = wrapped() + + assert result == "success" + mock_tracer.start_as_current_span.assert_called_once_with( + "google.api_core/TargetWithoutName", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.api_core", + "rpc.method": "TargetWithoutName", + }, + ) + + +def test_gapic_callable_otel_tracing_fallback_call_default(monkeypatch): + """Proves that _GapicCallable defaults method to 'call' if target has no __name__.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + + class NoNameTarget: + def __call__(self, *args, **kwargs): + return "success" + + target = NoNameTarget() + callable_obj = google.api_core.gapic_v1.method._GapicCallable( + target, None, None, None + ) + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_trace = mock.Mock() + mock_trace.get_tracer.return_value = mock_tracer + mock_trace.SpanKind.CLIENT = "CLIENT" + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + result = callable_obj() + + assert result == "success" + mock_tracer.start_as_current_span.assert_called_once_with( + "google.api_core/call", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.api_core", + "rpc.method": "call", + }, + ) + + +def test_wrap_method_otel_tracing_import_error(monkeypatch): + """Proves that if opentelemetry raises ImportError, execution proceeds gracefully.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="success") + + wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": None, + "opentelemetry.trace": None, + }, + ), + ): + result = wrapped() + + assert result == "success" + mock_target.assert_called_once() From b70b31ec6b67f338c0514f9bcb5894ad728b4d1e Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Tue, 8 Sep 2026 06:22:46 -0400 Subject: [PATCH 04/11] refactor(gapic): extract _extract_rpc_identity and add method_name to _GapicCallable --- .../google/api_core/gapic_v1/method.py | 60 ++++++++++++++----- .../google/api_core/gapic_v1/method_async.py | 2 + 2 files changed, 46 insertions(+), 16 deletions(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index dceb094f5568..4e3c4cb4dd1b 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -20,7 +20,7 @@ import enum import functools -from typing import List, Tuple +from typing import Any, List, Optional, Sequence, Tuple from google.api_core import _observability, grpc_helpers from google.api_core.gapic_v1 import client_info @@ -104,6 +104,36 @@ def _extract_metrics_header(metadata) -> Tuple[str, List[Tuple[str, str]]]: return metric_str, arbitrary_metadata +def _extract_rpc_identity( + target: Any, method_name: Optional[str] = None +) -> Tuple[str, str, str]: + """Extract (full_rpc_name, service_name, rpc_method_name) from an explicit method name or target callable. + + Args: + target: The underlying callable method. + method_name: Optional explicit RPC name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). + + Returns: + Tuple[str, str, str]: A 3-tuple of (full_rpc_name, service_name, rpc_method_name). + """ + if method_name: + method_str = method_name.lstrip("/") + service, _, method = method_str.rpartition("/") + return method_str, service, method + + raw_method = getattr(target, "_method", None) + if raw_method and isinstance(raw_method, (str, bytes)): + if isinstance(raw_method, bytes): + raw_method = raw_method.decode("utf-8") + method_str = raw_method.lstrip("/") + service, _, method = method_str.rpartition("/") + return method_str, service, method + + service = "google.api_core" + method = getattr(target, "__name__", "call") + return f"{service}/{method}", service, method + + class _GapicCallable(object): """Callable that applies retry, timeout, and metadata logic. @@ -123,6 +153,8 @@ class _GapicCallable(object): provided to the RPC method on every invocation. This is merged with any metadata specified during invocation. If ``None``, no additional metadata will be passed to the RPC method. + method_name (Optional[str]): The optional explicit full RPC method name + (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). """ def __init__( @@ -132,11 +164,15 @@ def __init__( timeout, compression, metadata=None, + method_name=None, ): self._target = target self._retry = retry self._timeout = timeout self._compression = compression + self._rpc_method_name, self._rpc_service, self._rpc_method = ( + _extract_rpc_identity(target, method_name) + ) # Pre-extract the x-goog-api-client header from the initialized metadata. self._x_goog_api_client, remaining = _extract_metrics_header(metadata) self._static_metadata = tuple(remaining) @@ -191,25 +227,13 @@ def __call__( from opentelemetry import trace tracer = trace.get_tracer("google.api_core") - raw_method = getattr(self._target, "_method", None) - if raw_method and isinstance(raw_method, (str, bytes)): - if isinstance(raw_method, bytes): - raw_method = raw_method.decode("utf-8") - method_str = raw_method.lstrip("/") - service, _, method = method_str.rpartition("/") - span_name = method_str - else: - service = "google.api_core" - method = getattr(self._target, "__name__", "call") - span_name = f"{service}/{method}" - with tracer.start_as_current_span( - span_name, + self._rpc_method_name, kind=trace.SpanKind.CLIENT, attributes={ "rpc.system": "grpc", - "rpc.service": service, - "rpc.method": method, + "rpc.service": self._rpc_service, + "rpc.method": self._rpc_method, }, ) as span: try: @@ -233,6 +257,7 @@ def wrap_method( client_info=client_info.DEFAULT_CLIENT_INFO, *, with_call=False, + method_name=None, ): """Wrap an RPC method with common behavior. @@ -316,6 +341,8 @@ def get_topic(name, timeout=None): return a tuple of (response, grpc.Call) instead of just the response. This is useful for extracting trailing metadata from unary calls. Defaults to False. + method_name (Optional[str]): Optional explicit full RPC method name + (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). Returns: Callable: A new callable that takes optional ``retry``, ``timeout``, @@ -343,5 +370,6 @@ def get_topic(name, timeout=None): default_timeout, default_compression, metadata=user_agent_metadata, + method_name=method_name, ) ) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method_async.py b/packages/google-api-core/google/api_core/gapic_v1/method_async.py index d361bf9f961f..1dcc45008e4f 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method_async.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method_async.py @@ -37,6 +37,7 @@ def wrap_method( default_compression=None, client_info=client_info.DEFAULT_CLIENT_INFO, kind=_DEFAULT_ASYNC_TRANSPORT_KIND, + method_name=None, ): """Wrap an async RPC method with common behavior. @@ -57,5 +58,6 @@ def wrap_method( default_timeout, default_compression, metadata=metadata, + method_name=method_name, ) ) From ae2a0da3a85267ad81a782c714c4ea7b4e897149 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Tue, 8 Sep 2026 06:23:18 -0400 Subject: [PATCH 05/11] refactor(gapic): support custom tracer_provider and cache tracer in _GapicCallable.__init__ --- .../google/api_core/gapic_v1/method.py | 25 ++++++++++++++++--- .../google/api_core/gapic_v1/method_async.py | 2 ++ 2 files changed, 24 insertions(+), 3 deletions(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index 4e3c4cb4dd1b..c89c9a16e6a3 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -155,6 +155,8 @@ class _GapicCallable(object): additional metadata will be passed to the RPC method. method_name (Optional[str]): The optional explicit full RPC method name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). + tracer_provider (Optional[Any]): Optional custom OpenTelemetry TracerProvider + to obtain the tracer from. """ def __init__( @@ -165,6 +167,7 @@ def __init__( compression, metadata=None, method_name=None, + tracer_provider=None, ): self._target = target self._retry = retry @@ -184,6 +187,19 @@ def __init__( else: self._default_metadata = self._static_metadata + # Resolve and cache the OpenTelemetry tracer once at initialization. + self._tracer = None + if _observability.is_otel_capabilities_enabled(): + try: + from opentelemetry import trace + + if tracer_provider is not None: + self._tracer = tracer_provider.get_tracer("google.api_core") + else: + self._tracer = trace.get_tracer("google.api_core") + except Exception: # pragma: NO COVER + self._tracer = None + def __call__( self, *args, timeout=DEFAULT, retry=DEFAULT, compression=DEFAULT, **kwargs ): @@ -222,12 +238,11 @@ def __call__( if self._compression is not None: kwargs["compression"] = compression - if _observability.is_otel_capabilities_enabled(): + if self._tracer is not None: try: from opentelemetry import trace - tracer = trace.get_tracer("google.api_core") - with tracer.start_as_current_span( + with self._tracer.start_as_current_span( self._rpc_method_name, kind=trace.SpanKind.CLIENT, attributes={ @@ -258,6 +273,7 @@ def wrap_method( *, with_call=False, method_name=None, + tracer_provider=None, ): """Wrap an RPC method with common behavior. @@ -343,6 +359,8 @@ def get_topic(name, timeout=None): Defaults to False. method_name (Optional[str]): Optional explicit full RPC method name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). + tracer_provider (Optional[Any]): Optional custom OpenTelemetry TracerProvider + to obtain the tracer from. Returns: Callable: A new callable that takes optional ``retry``, ``timeout``, @@ -371,5 +389,6 @@ def get_topic(name, timeout=None): default_compression, metadata=user_agent_metadata, method_name=method_name, + tracer_provider=tracer_provider, ) ) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method_async.py b/packages/google-api-core/google/api_core/gapic_v1/method_async.py index 1dcc45008e4f..f4968e6a5659 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method_async.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method_async.py @@ -38,6 +38,7 @@ def wrap_method( client_info=client_info.DEFAULT_CLIENT_INFO, kind=_DEFAULT_ASYNC_TRANSPORT_KIND, method_name=None, + tracer_provider=None, ): """Wrap an async RPC method with common behavior. @@ -59,5 +60,6 @@ def wrap_method( default_compression, metadata=metadata, method_name=method_name, + tracer_provider=tracer_provider, ) ) From 50d91d73ced5dececf0972e8394b5bfd2abed1b7 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Tue, 8 Sep 2026 06:23:46 -0400 Subject: [PATCH 06/11] feat(gapic): gate T3 method spans for HTTP transports and streaming RPCs --- .../google/api_core/gapic_v1/method.py | 22 ++++++++++++++++++- .../google/api_core/gapic_v1/method_async.py | 4 ++++ 2 files changed, 25 insertions(+), 1 deletion(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index c89c9a16e6a3..7bff5c944ea1 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -157,6 +157,10 @@ class _GapicCallable(object): (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). tracer_provider (Optional[Any]): Optional custom OpenTelemetry TracerProvider to obtain the tracer from. + rpc_system (Optional[str]): The RPC system (defaults to "grpc"). If not "grpc", + tracing will not be enabled. + is_streaming (bool): Whether the callable is a streaming RPC. Streaming RPC + tracing is currently gated and will not produce T3 spans. """ def __init__( @@ -168,11 +172,15 @@ def __init__( metadata=None, method_name=None, tracer_provider=None, + rpc_system="grpc", + is_streaming=False, ): self._target = target self._retry = retry self._timeout = timeout self._compression = compression + self._rpc_system = rpc_system + self._is_streaming = is_streaming self._rpc_method_name, self._rpc_service, self._rpc_method = ( _extract_rpc_identity(target, method_name) ) @@ -188,8 +196,13 @@ def __init__( self._default_metadata = self._static_metadata # Resolve and cache the OpenTelemetry tracer once at initialization. + # Tracing is gated to non-streaming gRPC calls for Tier 3 method spans. self._tracer = None - if _observability.is_otel_capabilities_enabled(): + if ( + rpc_system == "grpc" + and not is_streaming + and _observability.is_otel_capabilities_enabled() + ): try: from opentelemetry import trace @@ -274,6 +287,8 @@ def wrap_method( with_call=False, method_name=None, tracer_provider=None, + rpc_system="grpc", + is_streaming=False, ): """Wrap an RPC method with common behavior. @@ -361,6 +376,9 @@ def get_topic(name, timeout=None): (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). tracer_provider (Optional[Any]): Optional custom OpenTelemetry TracerProvider to obtain the tracer from. + rpc_system (Optional[str]): The RPC system (defaults to "grpc"). If not "grpc", + tracing will not be enabled. + is_streaming (bool): Whether the callable is a streaming RPC. Defaults to False. Returns: Callable: A new callable that takes optional ``retry``, ``timeout``, @@ -390,5 +408,7 @@ def get_topic(name, timeout=None): metadata=user_agent_metadata, method_name=method_name, tracer_provider=tracer_provider, + rpc_system=rpc_system, + is_streaming=is_streaming, ) ) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method_async.py b/packages/google-api-core/google/api_core/gapic_v1/method_async.py index f4968e6a5659..59d5763b9471 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method_async.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method_async.py @@ -39,6 +39,8 @@ def wrap_method( kind=_DEFAULT_ASYNC_TRANSPORT_KIND, method_name=None, tracer_provider=None, + rpc_system="grpc", + is_streaming=False, ): """Wrap an async RPC method with common behavior. @@ -61,5 +63,7 @@ def wrap_method( metadata=metadata, method_name=method_name, tracer_provider=tracer_provider, + rpc_system=rpc_system, + is_streaming=is_streaming, ) ) From fe04d0260715f378acc2085414c8aab539dcbd70 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Tue, 8 Sep 2026 06:24:14 -0400 Subject: [PATCH 07/11] refactor(gapic): use span_context_manager to guarantee single wrapped function invocation --- .../google/api_core/gapic_v1/method.py | 55 +++++++++++-------- 1 file changed, 33 insertions(+), 22 deletions(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index 7bff5c944ea1..cd3739f03ace 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -18,9 +18,10 @@ compression, pagination, and long-running operations to gRPC methods. """ +import contextlib import enum import functools -from typing import Any, List, Optional, Sequence, Tuple +from typing import Any, List, Optional, Tuple from google.api_core import _observability, grpc_helpers from google.api_core.gapic_v1 import client_info @@ -195,9 +196,11 @@ def __init__( else: self._default_metadata = self._static_metadata - # Resolve and cache the OpenTelemetry tracer once at initialization. + # Resolve and cache the OpenTelemetry tracer and attributes once at initialization. # Tracing is gated to non-streaming gRPC calls for Tier 3 method spans. self._tracer = None + self._span_name = None + self._span_attributes = None if ( rpc_system == "grpc" and not is_streaming @@ -210,8 +213,17 @@ def __init__( self._tracer = tracer_provider.get_tracer("google.api_core") else: self._tracer = trace.get_tracer("google.api_core") + + self._span_name = self._rpc_method_name + self._span_attributes = { + "rpc.system": "grpc", + "rpc.service": self._rpc_service, + "rpc.method": self._rpc_method, + } except Exception: # pragma: NO COVER self._tracer = None + self._span_name = None + self._span_attributes = None def __call__( self, *args, timeout=DEFAULT, retry=DEFAULT, compression=DEFAULT, **kwargs @@ -251,30 +263,29 @@ def __call__( if self._compression is not None: kwargs["compression"] = compression - if self._tracer is not None: + span_context_manager = contextlib.nullcontext() + if self._tracer is not None and self._span_name is not None: try: from opentelemetry import trace - with self._tracer.start_as_current_span( - self._rpc_method_name, + span_context_manager = self._tracer.start_as_current_span( + self._span_name, kind=trace.SpanKind.CLIENT, - attributes={ - "rpc.system": "grpc", - "rpc.service": self._rpc_service, - "rpc.method": self._rpc_method, - }, - ) as span: - try: - return wrapped_func(*args, **kwargs) - except Exception as exc: - span.record_exception(exc) - span.set_status(trace.StatusCode.ERROR, str(exc)) - raise - # If OpenTelemetry cannot be imported in the current environment, continue without tracing. - except ImportError: # pragma: NO COVER - pass - - return wrapped_func(*args, **kwargs) + attributes=self._span_attributes, + ) + except Exception: # pragma: NO COVER + span_context_manager = contextlib.nullcontext() + + with span_context_manager as span: + try: + return wrapped_func(*args, **kwargs) + except Exception as exc: + if span is not None and hasattr(span, "record_exception"): + from opentelemetry import trace + + span.record_exception(exc) + span.set_status(trace.StatusCode.ERROR, str(exc)) + raise def wrap_method( From 32dd79a101dba3784f6709e26f0b4a88079f8e65 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Tue, 8 Sep 2026 06:31:31 -0400 Subject: [PATCH 08/11] test(gapic): update and expand unit tests for T3 method tracing refactoring --- .../tests/asyncio/gapic/test_method_async.py | 47 ++++ .../tests/unit/gapic/test_method.py | 232 ++++++++++++++++-- 2 files changed, 261 insertions(+), 18 deletions(-) diff --git a/packages/google-api-core/tests/asyncio/gapic/test_method_async.py b/packages/google-api-core/tests/asyncio/gapic/test_method_async.py index e410acbdfaab..e27ae5854424 100644 --- a/packages/google-api-core/tests/asyncio/gapic/test_method_async.py +++ b/packages/google-api-core/tests/asyncio/gapic/test_method_async.py @@ -274,3 +274,50 @@ async def test_wrap_method_without_wrap_errors(): await wrapped_method() method.assert_not_called() + + +@pytest.mark.asyncio +async def test_wrap_method_async_with_otel_tracing(monkeypatch): + import sys + + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + fake_call = grpc_helpers_async.FakeUnaryUnaryCall(42) + method = mock.Mock(spec=aio.UnaryUnaryMultiCallable, return_value=fake_call) + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_trace = mock.Mock() + mock_trace.get_tracer.return_value = mock_tracer + mock_trace.SpanKind.CLIENT = "CLIENT" + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + wrapped_method = gapic_v1.method_async.wrap_method( + method, + method_name="google.test.AsyncService/AsyncMethod", + ) + result = await wrapped_method(1, 2, meep="moop") + + assert result == 42 + mock_tracer.start_as_current_span.assert_called_once_with( + "google.test.AsyncService/AsyncMethod", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.test.AsyncService", + "rpc.method": "AsyncMethod", + }, + ) diff --git a/packages/google-api-core/tests/unit/gapic/test_method.py b/packages/google-api-core/tests/unit/gapic/test_method.py index c337a31fe6f6..2a5d8f7f95a1 100644 --- a/packages/google-api-core/tests/unit/gapic/test_method.py +++ b/packages/google-api-core/tests/unit/gapic/test_method.py @@ -353,12 +353,12 @@ def test_wrap_method_otel_tracing_disabled(monkeypatch): """Proves that when OpenTelemetry tracing is disabled, no span is created.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "false") mock_target = mock.Mock(return_value="success") - wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) with mock.patch( "google.api_core._observability.is_otel_capabilities_enabled", return_value=False, ): + wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) assert wrapped() == "success" mock_target.assert_called_once() @@ -379,8 +379,6 @@ def test_wrap_method_otel_tracing_enabled_success(monkeypatch): mock_trace.get_tracer.return_value = mock_tracer mock_trace.SpanKind.CLIENT = "CLIENT" - wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) - with ( mock.patch( "google.api_core._observability.is_otel_capabilities_enabled", @@ -394,6 +392,7 @@ def test_wrap_method_otel_tracing_enabled_success(monkeypatch): }, ), ): + wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) result = wrapped() assert result == "success" @@ -408,6 +407,160 @@ def test_wrap_method_otel_tracing_enabled_success(monkeypatch): ) +def test_wrap_method_otel_tracing_explicit_method_name(monkeypatch): + """Proves that passing method_name explicitly configures span name and attributes without introspection.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="success") + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_trace = mock.Mock() + mock_trace.get_tracer.return_value = mock_tracer + mock_trace.SpanKind.CLIENT = "CLIENT" + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, + method_name="/google.cloud.secretmanager.v1.SecretManagerService/CreateSecret", + ) + result = wrapped() + + assert result == "success" + mock_tracer.start_as_current_span.assert_called_once_with( + "google.cloud.secretmanager.v1.SecretManagerService/CreateSecret", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.cloud.secretmanager.v1.SecretManagerService", + "rpc.method": "CreateSecret", + }, + ) + + +def test_wrap_method_otel_tracing_custom_tracer_provider(monkeypatch): + """Proves that providing a custom tracer_provider uses that provider to obtain the tracer.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="success") + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_provider = mock.Mock() + mock_provider.get_tracer.return_value = mock_tracer + + mock_trace = mock.Mock() + mock_trace.SpanKind.CLIENT = "CLIENT" + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, + method_name="google.test.Service/TestMethod", + tracer_provider=mock_provider, + ) + result = wrapped() + + assert result == "success" + mock_provider.get_tracer.assert_called_once_with("google.api_core") + mock_tracer.start_as_current_span.assert_called_once_with( + "google.test.Service/TestMethod", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.test.Service", + "rpc.method": "TestMethod", + }, + ) + + +def test_wrap_method_otel_tracing_gated_http(monkeypatch): + """Proves that HTTP transports (rpc_system != 'grpc') do not generate T3 spans.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="success") + + mock_trace = mock.Mock() + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, + method_name="google.test.Service/HttpCall", + rpc_system="http", + ) + result = wrapped() + + assert result == "success" + mock_trace.get_tracer.assert_not_called() + + +def test_wrap_method_otel_tracing_gated_streaming(monkeypatch): + """Proves that streaming RPCs (is_streaming=True) do not generate T3 spans.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="success") + + mock_trace = mock.Mock() + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, + method_name="google.test.Service/StreamCall", + is_streaming=True, + ) + result = wrapped() + + assert result == "success" + mock_trace.get_tracer.assert_not_called() + + def test_wrap_method_otel_tracing_enabled_error(monkeypatch): """Proves that when an RPC fails, the T3 client span records the exception and error status.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") @@ -426,8 +579,6 @@ def test_wrap_method_otel_tracing_enabled_error(monkeypatch): mock_trace.SpanKind.CLIENT = "CLIENT" mock_trace.StatusCode.ERROR = "ERROR" - wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) - with ( mock.patch( "google.api_core._observability.is_otel_capabilities_enabled", @@ -441,9 +592,11 @@ def test_wrap_method_otel_tracing_enabled_error(monkeypatch): }, ), ): + wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) with pytest.raises(RuntimeError): wrapped() + mock_target.assert_called_once() mock_span.record_exception.assert_called_once_with(err) mock_span.set_status.assert_called_once_with("ERROR", str(err)) @@ -464,10 +617,6 @@ def test_wrap_method_otel_tracing_bytes_method(monkeypatch): mock_trace.get_tracer.return_value = mock_tracer mock_trace.SpanKind.CLIENT = "CLIENT" - wrapped = google.api_core.gapic_v1.method.wrap_method( - mock_target, default_timeout=60 - ) - with ( mock.patch( "google.api_core._observability.is_otel_capabilities_enabled", @@ -481,6 +630,9 @@ def test_wrap_method_otel_tracing_bytes_method(monkeypatch): }, ), ): + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, default_timeout=60 + ) result = wrapped() assert result == "success" @@ -510,8 +662,6 @@ def custom_rpc(*args, **kwargs): mock_trace.get_tracer.return_value = mock_tracer mock_trace.SpanKind.CLIENT = "CLIENT" - wrapped = google.api_core.gapic_v1.method.wrap_method(custom_rpc) - with ( mock.patch( "google.api_core._observability.is_otel_capabilities_enabled", @@ -525,6 +675,7 @@ def custom_rpc(*args, **kwargs): }, ), ): + wrapped = google.api_core.gapic_v1.method.wrap_method(custom_rpc) result = wrapped() assert result == "success" @@ -557,8 +708,6 @@ def __call__(self, *args, **kwargs): mock_trace.get_tracer.return_value = mock_tracer mock_trace.SpanKind.CLIENT = "CLIENT" - wrapped = google.api_core.gapic_v1.method.wrap_method(TargetWithoutName()) - with ( mock.patch( "google.api_core._observability.is_otel_capabilities_enabled", @@ -572,6 +721,7 @@ def __call__(self, *args, **kwargs): }, ), ): + wrapped = google.api_core.gapic_v1.method.wrap_method(TargetWithoutName()) result = wrapped() assert result == "success" @@ -595,9 +745,6 @@ def __call__(self, *args, **kwargs): return "success" target = NoNameTarget() - callable_obj = google.api_core.gapic_v1.method._GapicCallable( - target, None, None, None - ) mock_span = mock.MagicMock() mock_tracer = mock.MagicMock() @@ -620,6 +767,9 @@ def __call__(self, *args, **kwargs): }, ), ): + callable_obj = google.api_core.gapic_v1.method._GapicCallable( + target, None, None, None + ) result = callable_obj() assert result == "success" @@ -639,8 +789,6 @@ def test_wrap_method_otel_tracing_import_error(monkeypatch): monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") mock_target = mock.Mock(return_value="success") - wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) - with ( mock.patch( "google.api_core._observability.is_otel_capabilities_enabled", @@ -654,7 +802,55 @@ def test_wrap_method_otel_tracing_import_error(monkeypatch): }, ), ): + wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) result = wrapped() assert result == "success" mock_target.assert_called_once() + + +def test_wrap_method_async_otel_tracing(monkeypatch): + """Proves that method_async.wrap_method correctly passes OTel arguments to _GapicCallable.""" + from google.api_core.gapic_v1 import method_async + + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="async_success") + + mock_span = mock.MagicMock() + mock_tracer = mock.MagicMock() + mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + + mock_trace = mock.Mock() + mock_trace.get_tracer.return_value = mock_tracer + mock_trace.SpanKind.CLIENT = "CLIENT" + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + wrapped = method_async.wrap_method( + mock_target, + kind=None, + method_name="google.test.AsyncService/AsyncMethod", + ) + result = wrapped() + + assert result == "async_success" + mock_tracer.start_as_current_span.assert_called_once_with( + "google.test.AsyncService/AsyncMethod", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.test.AsyncService", + "rpc.method": "AsyncMethod", + }, + ) From 420582ff76494233bbf9c4728c8de0163d2bdde8 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Tue, 8 Sep 2026 07:57:34 -0400 Subject: [PATCH 09/11] refactor(gapic): simplify wrap_method with client_options and method_name Streamline wrap_method and _GapicCallable signatures by replacing loose primitive arguments (method_name, tracer_provider, rpc_system, is_streaming) with client_options and method_name. - Pass client_options to extract tracer_provider if configured, falling back to the global OpenTelemetry tracer provider. - Gate Tier 3 method span creation naturally on whether method_name is provided. When method_name is omitted (e.g. For streaming RPCs or uninstrumented transports), span generation is bypassed with zero overhead. - Update unit tests to validate client_options and method_name gating, and maintain 100% statement and branch coverage. --- .../google/api_core/gapic_v1/method.py | 88 +++--- .../google/api_core/gapic_v1/method_async.py | 9 +- .../tests/unit/gapic/test_method.py | 278 +++--------------- 3 files changed, 76 insertions(+), 299 deletions(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index cd3739f03ace..776aef012533 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -21,7 +21,7 @@ import contextlib import enum import functools -from typing import Any, List, Optional, Tuple +from typing import List, Tuple from google.api_core import _observability, grpc_helpers from google.api_core.gapic_v1 import client_info @@ -106,33 +106,21 @@ def _extract_metrics_header(metadata) -> Tuple[str, List[Tuple[str, str]]]: def _extract_rpc_identity( - target: Any, method_name: Optional[str] = None + method_name: str, ) -> Tuple[str, str, str]: - """Extract (full_rpc_name, service_name, rpc_method_name) from an explicit method name or target callable. + """Extract (full_rpc_name, service_name, rpc_method_name) from an explicit method name. Args: - target: The underlying callable method. - method_name: Optional explicit RPC name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). + method_name: Explicit RPC name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). Returns: Tuple[str, str, str]: A 3-tuple of (full_rpc_name, service_name, rpc_method_name). """ - if method_name: - method_str = method_name.lstrip("/") - service, _, method = method_str.rpartition("/") - return method_str, service, method - - raw_method = getattr(target, "_method", None) - if raw_method and isinstance(raw_method, (str, bytes)): - if isinstance(raw_method, bytes): - raw_method = raw_method.decode("utf-8") - method_str = raw_method.lstrip("/") - service, _, method = method_str.rpartition("/") - return method_str, service, method - - service = "google.api_core" - method = getattr(target, "__name__", "call") - return f"{service}/{method}", service, method + if isinstance(method_name, bytes): + method_name = method_name.decode("utf-8") + method_str = method_name.lstrip("/") + service, _, method = method_str.rpartition("/") + return method_str, service, method class _GapicCallable(object): @@ -154,14 +142,13 @@ class _GapicCallable(object): provided to the RPC method on every invocation. This is merged with any metadata specified during invocation. If ``None``, no additional metadata will be passed to the RPC method. + client_options + (Optional[google.api_core.client_options.ClientOptions]): + Client options used to configure client-level behavior, such as + custom OpenTelemetry tracer providers. Defaults to None. method_name (Optional[str]): The optional explicit full RPC method name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). - tracer_provider (Optional[Any]): Optional custom OpenTelemetry TracerProvider - to obtain the tracer from. - rpc_system (Optional[str]): The RPC system (defaults to "grpc"). If not "grpc", - tracing will not be enabled. - is_streaming (bool): Whether the callable is a streaming RPC. Streaming RPC - tracing is currently gated and will not produce T3 spans. + If omitted or None, method-level tracing spans are not generated. """ def __init__( @@ -171,20 +158,16 @@ def __init__( timeout, compression, metadata=None, + client_options=None, method_name=None, - tracer_provider=None, - rpc_system="grpc", - is_streaming=False, ): self._target = target self._retry = retry self._timeout = timeout self._compression = compression - self._rpc_system = rpc_system - self._is_streaming = is_streaming - self._rpc_method_name, self._rpc_service, self._rpc_method = ( - _extract_rpc_identity(target, method_name) - ) + self._client_options = client_options + self._method_name = method_name + # Pre-extract the x-goog-api-client header from the initialized metadata. self._x_goog_api_client, remaining = _extract_metrics_header(metadata) self._static_metadata = tuple(remaining) @@ -197,24 +180,27 @@ def __init__( self._default_metadata = self._static_metadata # Resolve and cache the OpenTelemetry tracer and attributes once at initialization. - # Tracing is gated to non-streaming gRPC calls for Tier 3 method spans. + # Tracing is gated to calls where an explicit method_name is provided. self._tracer = None self._span_name = None self._span_attributes = None - if ( - rpc_system == "grpc" - and not is_streaming - and _observability.is_otel_capabilities_enabled() - ): + if method_name is not None and _observability.is_otel_capabilities_enabled(): try: from opentelemetry import trace + tracer_provider = ( + getattr(client_options, "tracer_provider", None) + if client_options is not None + else None + ) if tracer_provider is not None: self._tracer = tracer_provider.get_tracer("google.api_core") else: self._tracer = trace.get_tracer("google.api_core") - self._span_name = self._rpc_method_name + self._span_name, self._rpc_service, self._rpc_method = ( + _extract_rpc_identity(method_name) + ) self._span_attributes = { "rpc.system": "grpc", "rpc.service": self._rpc_service, @@ -296,10 +282,8 @@ def wrap_method( client_info=client_info.DEFAULT_CLIENT_INFO, *, with_call=False, + client_options=None, method_name=None, - tracer_provider=None, - rpc_system="grpc", - is_streaming=False, ): """Wrap an RPC method with common behavior. @@ -383,13 +367,13 @@ def get_topic(name, timeout=None): return a tuple of (response, grpc.Call) instead of just the response. This is useful for extracting trailing metadata from unary calls. Defaults to False. + client_options + (Optional[google.api_core.client_options.ClientOptions]): + Client options used to configure client-level behavior, such as + custom OpenTelemetry tracer providers. Defaults to None. method_name (Optional[str]): Optional explicit full RPC method name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). - tracer_provider (Optional[Any]): Optional custom OpenTelemetry TracerProvider - to obtain the tracer from. - rpc_system (Optional[str]): The RPC system (defaults to "grpc"). If not "grpc", - tracing will not be enabled. - is_streaming (bool): Whether the callable is a streaming RPC. Defaults to False. + If omitted or None, method-level tracing spans are not generated. Returns: Callable: A new callable that takes optional ``retry``, ``timeout``, @@ -417,9 +401,7 @@ def get_topic(name, timeout=None): default_timeout, default_compression, metadata=user_agent_metadata, + client_options=client_options, method_name=method_name, - tracer_provider=tracer_provider, - rpc_system=rpc_system, - is_streaming=is_streaming, ) ) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method_async.py b/packages/google-api-core/google/api_core/gapic_v1/method_async.py index 59d5763b9471..752988cb3f72 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method_async.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method_async.py @@ -37,10 +37,9 @@ def wrap_method( default_compression=None, client_info=client_info.DEFAULT_CLIENT_INFO, kind=_DEFAULT_ASYNC_TRANSPORT_KIND, + *, + client_options=None, method_name=None, - tracer_provider=None, - rpc_system="grpc", - is_streaming=False, ): """Wrap an async RPC method with common behavior. @@ -61,9 +60,7 @@ def wrap_method( default_timeout, default_compression, metadata=metadata, + client_options=client_options, method_name=method_name, - tracer_provider=tracer_provider, - rpc_system=rpc_system, - is_streaming=is_streaming, ) ) diff --git a/packages/google-api-core/tests/unit/gapic/test_method.py b/packages/google-api-core/tests/unit/gapic/test_method.py index 2a5d8f7f95a1..250de22bd60c 100644 --- a/packages/google-api-core/tests/unit/gapic/test_method.py +++ b/packages/google-api-core/tests/unit/gapic/test_method.py @@ -27,6 +27,7 @@ import google.api_core.gapic_v1.client_info import google.api_core.gapic_v1.method import google.api_core.page_iterator +from google.api_core import client_options as client_options_lib from google.api_core import exceptions, retry, timeout @@ -358,26 +359,20 @@ def test_wrap_method_otel_tracing_disabled(monkeypatch): "google.api_core._observability.is_otel_capabilities_enabled", return_value=False, ): - wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, + method_name="google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", + ) assert wrapped() == "success" mock_target.assert_called_once() -def test_wrap_method_otel_tracing_enabled_success(monkeypatch): - """Proves that when OpenTelemetry tracing is enabled, a T3 client span is started.""" +def test_wrap_method_otel_tracing_omitted_method_name_skips_span(monkeypatch): + """Proves that when method_name is omitted (e.g. streaming or uninstrumented), no span is created.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") mock_target = mock.Mock(return_value="success") - mock_target._method = ( - "/google.cloud.secretmanager.v1.SecretManagerService/ListSecrets" - ) - - mock_span = mock.MagicMock() - mock_tracer = mock.MagicMock() - mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span mock_trace = mock.Mock() - mock_trace.get_tracer.return_value = mock_tracer - mock_trace.SpanKind.CLIENT = "CLIENT" with ( mock.patch( @@ -396,19 +391,11 @@ def test_wrap_method_otel_tracing_enabled_success(monkeypatch): result = wrapped() assert result == "success" - mock_tracer.start_as_current_span.assert_called_once_with( - "google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", - kind="CLIENT", - attributes={ - "rpc.system": "grpc", - "rpc.service": "google.cloud.secretmanager.v1.SecretManagerService", - "rpc.method": "ListSecrets", - }, - ) + mock_trace.get_tracer.assert_not_called() -def test_wrap_method_otel_tracing_explicit_method_name(monkeypatch): - """Proves that passing method_name explicitly configures span name and attributes without introspection.""" +def test_wrap_method_otel_tracing_enabled_success(monkeypatch): + """Proves that when OpenTelemetry tracing is enabled and method_name is passed, a T3 client span is started.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") mock_target = mock.Mock(return_value="success") @@ -435,24 +422,24 @@ def test_wrap_method_otel_tracing_explicit_method_name(monkeypatch): ): wrapped = google.api_core.gapic_v1.method.wrap_method( mock_target, - method_name="/google.cloud.secretmanager.v1.SecretManagerService/CreateSecret", + method_name="/google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", ) result = wrapped() assert result == "success" mock_tracer.start_as_current_span.assert_called_once_with( - "google.cloud.secretmanager.v1.SecretManagerService/CreateSecret", + "google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", kind="CLIENT", attributes={ "rpc.system": "grpc", "rpc.service": "google.cloud.secretmanager.v1.SecretManagerService", - "rpc.method": "CreateSecret", + "rpc.method": "ListSecrets", }, ) -def test_wrap_method_otel_tracing_custom_tracer_provider(monkeypatch): - """Proves that providing a custom tracer_provider uses that provider to obtain the tracer.""" +def test_wrap_method_otel_tracing_custom_client_options(monkeypatch): + """Proves that providing client_options with a custom tracer_provider uses that provider.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") mock_target = mock.Mock(return_value="success") @@ -466,6 +453,8 @@ def test_wrap_method_otel_tracing_custom_tracer_provider(monkeypatch): mock_trace = mock.Mock() mock_trace.SpanKind.CLIENT = "CLIENT" + client_options = client_options_lib.ClientOptions(tracer_provider=mock_provider) + with ( mock.patch( "google.api_core._observability.is_otel_capabilities_enabled", @@ -481,8 +470,8 @@ def test_wrap_method_otel_tracing_custom_tracer_provider(monkeypatch): ): wrapped = google.api_core.gapic_v1.method.wrap_method( mock_target, + client_options=client_options, method_name="google.test.Service/TestMethod", - tracer_provider=mock_provider, ) result = wrapped() @@ -499,76 +488,11 @@ def test_wrap_method_otel_tracing_custom_tracer_provider(monkeypatch): ) -def test_wrap_method_otel_tracing_gated_http(monkeypatch): - """Proves that HTTP transports (rpc_system != 'grpc') do not generate T3 spans.""" - monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") - mock_target = mock.Mock(return_value="success") - - mock_trace = mock.Mock() - - with ( - mock.patch( - "google.api_core._observability.is_otel_capabilities_enabled", - return_value=True, - ), - mock.patch.dict( - sys.modules, - { - "opentelemetry": mock.Mock(trace=mock_trace), - "opentelemetry.trace": mock_trace, - }, - ), - ): - wrapped = google.api_core.gapic_v1.method.wrap_method( - mock_target, - method_name="google.test.Service/HttpCall", - rpc_system="http", - ) - result = wrapped() - - assert result == "success" - mock_trace.get_tracer.assert_not_called() - - -def test_wrap_method_otel_tracing_gated_streaming(monkeypatch): - """Proves that streaming RPCs (is_streaming=True) do not generate T3 spans.""" - monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") - mock_target = mock.Mock(return_value="success") - - mock_trace = mock.Mock() - - with ( - mock.patch( - "google.api_core._observability.is_otel_capabilities_enabled", - return_value=True, - ), - mock.patch.dict( - sys.modules, - { - "opentelemetry": mock.Mock(trace=mock_trace), - "opentelemetry.trace": mock_trace, - }, - ), - ): - wrapped = google.api_core.gapic_v1.method.wrap_method( - mock_target, - method_name="google.test.Service/StreamCall", - is_streaming=True, - ) - result = wrapped() - - assert result == "success" - mock_trace.get_tracer.assert_not_called() - - def test_wrap_method_otel_tracing_enabled_error(monkeypatch): """Proves that when an RPC fails, the T3 client span records the exception and error status.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") err = RuntimeError("gRPC connection reset") mock_target = mock.Mock(side_effect=err) - mock_target._method = ( - "/google.cloud.secretmanager.v1.SecretManagerService/ListSecrets" - ) mock_span = mock.MagicMock() mock_tracer = mock.MagicMock() @@ -592,7 +516,10 @@ def test_wrap_method_otel_tracing_enabled_error(monkeypatch): }, ), ): - wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, + method_name="/google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", + ) with pytest.raises(RuntimeError): wrapped() @@ -602,12 +529,9 @@ def test_wrap_method_otel_tracing_enabled_error(monkeypatch): def test_wrap_method_otel_tracing_bytes_method(monkeypatch): - """Proves that when raw _method is bytes, it is decoded properly to utf-8.""" + """Proves that when method_name is bytes, it is decoded properly to utf-8.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") mock_target = mock.Mock(return_value="success") - mock_target._method = ( - b"/google.cloud.secretmanager.v1.SecretManagerService/ListSecrets" - ) mock_span = mock.MagicMock() mock_tracer = mock.MagicMock() @@ -631,7 +555,9 @@ def test_wrap_method_otel_tracing_bytes_method(monkeypatch): ), ): wrapped = google.api_core.gapic_v1.method.wrap_method( - mock_target, default_timeout=60 + mock_target, + default_timeout=60, + method_name=b"/google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", ) result = wrapped() @@ -647,143 +573,6 @@ def test_wrap_method_otel_tracing_bytes_method(monkeypatch): ) -def test_wrap_method_otel_tracing_fallback_with_name(monkeypatch): - """Proves that when raw _method is absent, fallback uses target.__name__.""" - monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") - - def custom_rpc(*args, **kwargs): - return "success" - - mock_span = mock.MagicMock() - mock_tracer = mock.MagicMock() - mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span - - mock_trace = mock.Mock() - mock_trace.get_tracer.return_value = mock_tracer - mock_trace.SpanKind.CLIENT = "CLIENT" - - with ( - mock.patch( - "google.api_core._observability.is_otel_capabilities_enabled", - return_value=True, - ), - mock.patch.dict( - sys.modules, - { - "opentelemetry": mock.Mock(trace=mock_trace), - "opentelemetry.trace": mock_trace, - }, - ), - ): - wrapped = google.api_core.gapic_v1.method.wrap_method(custom_rpc) - result = wrapped() - - assert result == "success" - mock_tracer.start_as_current_span.assert_called_once_with( - "google.api_core/custom_rpc", - kind="CLIENT", - attributes={ - "rpc.system": "grpc", - "rpc.service": "google.api_core", - "rpc.method": "custom_rpc", - }, - ) - - -def test_wrap_method_otel_tracing_fallback_without_name(monkeypatch): - """Proves that when raw _method is absent and target has no explicit __name__, - fallback uses target class name assigned by error wrapper. - """ - monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") - - class TargetWithoutName: - def __call__(self, *args, **kwargs): - return "success" - - mock_span = mock.MagicMock() - mock_tracer = mock.MagicMock() - mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span - - mock_trace = mock.Mock() - mock_trace.get_tracer.return_value = mock_tracer - mock_trace.SpanKind.CLIENT = "CLIENT" - - with ( - mock.patch( - "google.api_core._observability.is_otel_capabilities_enabled", - return_value=True, - ), - mock.patch.dict( - sys.modules, - { - "opentelemetry": mock.Mock(trace=mock_trace), - "opentelemetry.trace": mock_trace, - }, - ), - ): - wrapped = google.api_core.gapic_v1.method.wrap_method(TargetWithoutName()) - result = wrapped() - - assert result == "success" - mock_tracer.start_as_current_span.assert_called_once_with( - "google.api_core/TargetWithoutName", - kind="CLIENT", - attributes={ - "rpc.system": "grpc", - "rpc.service": "google.api_core", - "rpc.method": "TargetWithoutName", - }, - ) - - -def test_gapic_callable_otel_tracing_fallback_call_default(monkeypatch): - """Proves that _GapicCallable defaults method to 'call' if target has no __name__.""" - monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") - - class NoNameTarget: - def __call__(self, *args, **kwargs): - return "success" - - target = NoNameTarget() - - mock_span = mock.MagicMock() - mock_tracer = mock.MagicMock() - mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span - - mock_trace = mock.Mock() - mock_trace.get_tracer.return_value = mock_tracer - mock_trace.SpanKind.CLIENT = "CLIENT" - - with ( - mock.patch( - "google.api_core._observability.is_otel_capabilities_enabled", - return_value=True, - ), - mock.patch.dict( - sys.modules, - { - "opentelemetry": mock.Mock(trace=mock_trace), - "opentelemetry.trace": mock_trace, - }, - ), - ): - callable_obj = google.api_core.gapic_v1.method._GapicCallable( - target, None, None, None - ) - result = callable_obj() - - assert result == "success" - mock_tracer.start_as_current_span.assert_called_once_with( - "google.api_core/call", - kind="CLIENT", - attributes={ - "rpc.system": "grpc", - "rpc.service": "google.api_core", - "rpc.method": "call", - }, - ) - - def test_wrap_method_otel_tracing_import_error(monkeypatch): """Proves that if opentelemetry raises ImportError, execution proceeds gracefully.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") @@ -802,7 +591,10 @@ def test_wrap_method_otel_tracing_import_error(monkeypatch): }, ), ): - wrapped = google.api_core.gapic_v1.method.wrap_method(mock_target) + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, + method_name="google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", + ) result = wrapped() assert result == "success" @@ -810,7 +602,7 @@ def test_wrap_method_otel_tracing_import_error(monkeypatch): def test_wrap_method_async_otel_tracing(monkeypatch): - """Proves that method_async.wrap_method correctly passes OTel arguments to _GapicCallable.""" + """Proves that method_async.wrap_method correctly passes client_options and method_name to _GapicCallable.""" from google.api_core.gapic_v1 import method_async monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") @@ -820,10 +612,14 @@ def test_wrap_method_async_otel_tracing(monkeypatch): mock_tracer = mock.MagicMock() mock_tracer.start_as_current_span.return_value.__enter__.return_value = mock_span + mock_provider = mock.Mock() + mock_provider.get_tracer.return_value = mock_tracer + mock_trace = mock.Mock() - mock_trace.get_tracer.return_value = mock_tracer mock_trace.SpanKind.CLIENT = "CLIENT" + client_options = client_options_lib.ClientOptions(tracer_provider=mock_provider) + with ( mock.patch( "google.api_core._observability.is_otel_capabilities_enabled", @@ -840,11 +636,13 @@ def test_wrap_method_async_otel_tracing(monkeypatch): wrapped = method_async.wrap_method( mock_target, kind=None, + client_options=client_options, method_name="google.test.AsyncService/AsyncMethod", ) result = wrapped() assert result == "async_success" + mock_provider.get_tracer.assert_called_once_with("google.api_core") mock_tracer.start_as_current_span.assert_called_once_with( "google.test.AsyncService/AsyncMethod", kind="CLIENT", From 2a520995e109105e50cf595c21d3c7493d0820f8 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Tue, 8 Sep 2026 08:18:08 -0400 Subject: [PATCH 10/11] feat(gapic): add explicit trace parameter to wrap_method Separate RPC method identity from tracing behavior by adding an explicit, keyword-only 'trace: bool = True' parameter to wrap_method and _GapicCallable. - Allows methods to retain their true method_name while explicitly disabling Tier 3 method span creation (e.g. For streaming calls or unsupported transports via trace=False). - Updates method_async.wrap_method to forward trace to _GapicCallable. - Adds sync and async unit tests verifying that trace=False bypasses span creation even when method_name is provided, maintaining 100% statement and branch coverage. --- .../google/api_core/gapic_v1/method.py | 22 +++++-- .../google/api_core/gapic_v1/method_async.py | 2 + .../tests/unit/gapic/test_method.py | 64 +++++++++++++++++++ 3 files changed, 84 insertions(+), 4 deletions(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index 776aef012533..3149b458231b 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -148,7 +148,10 @@ class _GapicCallable(object): custom OpenTelemetry tracer providers. Defaults to None. method_name (Optional[str]): The optional explicit full RPC method name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). - If omitted or None, method-level tracing spans are not generated. + Used to identify the RPC for observability and tracing. + trace (bool): Whether to create OpenTelemetry Tier 3 tracing spans for this + callable. Defaults to True. When False, or when method_name is None, + tracing spans are bypassed. """ def __init__( @@ -160,6 +163,7 @@ def __init__( metadata=None, client_options=None, method_name=None, + trace=True, ): self._target = target self._retry = retry @@ -167,6 +171,7 @@ def __init__( self._compression = compression self._client_options = client_options self._method_name = method_name + self._trace = trace # Pre-extract the x-goog-api-client header from the initialized metadata. self._x_goog_api_client, remaining = _extract_metrics_header(metadata) @@ -180,11 +185,15 @@ def __init__( self._default_metadata = self._static_metadata # Resolve and cache the OpenTelemetry tracer and attributes once at initialization. - # Tracing is gated to calls where an explicit method_name is provided. + # Tracing is gated to calls where trace is True and an explicit method_name is provided. self._tracer = None self._span_name = None self._span_attributes = None - if method_name is not None and _observability.is_otel_capabilities_enabled(): + if ( + trace + and method_name is not None + and _observability.is_otel_capabilities_enabled() + ): try: from opentelemetry import trace @@ -284,6 +293,7 @@ def wrap_method( with_call=False, client_options=None, method_name=None, + trace=True, ): """Wrap an RPC method with common behavior. @@ -373,7 +383,10 @@ def get_topic(name, timeout=None): custom OpenTelemetry tracer providers. Defaults to None. method_name (Optional[str]): Optional explicit full RPC method name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). - If omitted or None, method-level tracing spans are not generated. + Used to identify the RPC for observability and tracing. + trace (bool): Whether to create OpenTelemetry Tier 3 tracing spans for this + callable. Defaults to True. When False, or when method_name is None, + tracing spans are bypassed. Returns: Callable: A new callable that takes optional ``retry``, ``timeout``, @@ -403,5 +416,6 @@ def get_topic(name, timeout=None): metadata=user_agent_metadata, client_options=client_options, method_name=method_name, + trace=trace, ) ) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method_async.py b/packages/google-api-core/google/api_core/gapic_v1/method_async.py index 752988cb3f72..9ea110717352 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method_async.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method_async.py @@ -40,6 +40,7 @@ def wrap_method( *, client_options=None, method_name=None, + trace=True, ): """Wrap an async RPC method with common behavior. @@ -62,5 +63,6 @@ def wrap_method( metadata=metadata, client_options=client_options, method_name=method_name, + trace=trace, ) ) diff --git a/packages/google-api-core/tests/unit/gapic/test_method.py b/packages/google-api-core/tests/unit/gapic/test_method.py index 250de22bd60c..5f3330c8a689 100644 --- a/packages/google-api-core/tests/unit/gapic/test_method.py +++ b/packages/google-api-core/tests/unit/gapic/test_method.py @@ -394,6 +394,37 @@ def test_wrap_method_otel_tracing_omitted_method_name_skips_span(monkeypatch): mock_trace.get_tracer.assert_not_called() +def test_wrap_method_otel_tracing_explicit_trace_false_skips_span(monkeypatch): + """Proves that when trace=False is explicitly passed (e.g. streaming call), no span is created.""" + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="success") + + mock_trace = mock.Mock() + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, + method_name="/google.cloud.secretmanager.v1.SecretManagerService/StreamingRead", + trace=False, + ) + result = wrapped() + + assert result == "success" + mock_trace.get_tracer.assert_not_called() + + def test_wrap_method_otel_tracing_enabled_success(monkeypatch): """Proves that when OpenTelemetry tracing is enabled and method_name is passed, a T3 client span is started.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") @@ -652,3 +683,36 @@ def test_wrap_method_async_otel_tracing(monkeypatch): "rpc.method": "AsyncMethod", }, ) + + +def test_wrap_method_async_otel_tracing_trace_false_skips_span(monkeypatch): + """Proves that method_async.wrap_method with trace=False skips span creation.""" + from google.api_core.gapic_v1 import method_async + + monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") + mock_target = mock.Mock(return_value="async_success") + mock_trace = mock.Mock() + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": mock.Mock(trace=mock_trace), + "opentelemetry.trace": mock_trace, + }, + ), + ): + wrapped = method_async.wrap_method( + mock_target, + kind=None, + method_name="google.test.AsyncService/AsyncMethod", + trace=False, + ) + result = wrapped() + + assert result == "async_success" + mock_trace.get_tracer.assert_not_called() From f81e4be822cc1f9d648b2386805fa632ba7f4a4c Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Tue, 8 Sep 2026 09:00:14 -0400 Subject: [PATCH 11/11] refactor(gapic): replace trace with is_streaming in wrap_method Replace subsystem-specific 'trace' boolean with physical method topology 'is_streaming: bool = False' across wrap_method and _GapicCallable. - Encapsulates method characteristics cleanly without requiring callers/generators to act as observability policy engines. - Centralizes enablement checks in _observability.is_otel_capabilities_enabled(client_options), avoiding duplication of feature flag logic. - Prepares wrap_method for future metrics and logging additions without signature churn. - Updates unit tests to verify is_streaming=True gating for sync and async callables, maintaining 100% statement and branch coverage. --- .../google/api_core/gapic_v1/method.py | 28 +++++++++---------- .../google/api_core/gapic_v1/method_async.py | 4 +-- .../tests/unit/gapic/test_method.py | 12 ++++---- 3 files changed, 21 insertions(+), 23 deletions(-) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method.py b/packages/google-api-core/google/api_core/gapic_v1/method.py index 3149b458231b..e150bf5494d2 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method.py @@ -148,10 +148,9 @@ class _GapicCallable(object): custom OpenTelemetry tracer providers. Defaults to None. method_name (Optional[str]): The optional explicit full RPC method name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). - Used to identify the RPC for observability and tracing. - trace (bool): Whether to create OpenTelemetry Tier 3 tracing spans for this - callable. Defaults to True. When False, or when method_name is None, - tracing spans are bypassed. + Used to identify the RPC for observability. + is_streaming (bool): Whether the RPC method is streaming. Defaults to False. + Streaming methods are currently gated and do not generate Tier 3 spans. """ def __init__( @@ -163,7 +162,7 @@ def __init__( metadata=None, client_options=None, method_name=None, - trace=True, + is_streaming=False, ): self._target = target self._retry = retry @@ -171,7 +170,7 @@ def __init__( self._compression = compression self._client_options = client_options self._method_name = method_name - self._trace = trace + self._is_streaming = is_streaming # Pre-extract the x-goog-api-client header from the initialized metadata. self._x_goog_api_client, remaining = _extract_metrics_header(metadata) @@ -185,14 +184,14 @@ def __init__( self._default_metadata = self._static_metadata # Resolve and cache the OpenTelemetry tracer and attributes once at initialization. - # Tracing is gated to calls where trace is True and an explicit method_name is provided. + # Tracing is gated to non-streaming calls where an explicit method_name is provided. self._tracer = None self._span_name = None self._span_attributes = None if ( - trace + not is_streaming and method_name is not None - and _observability.is_otel_capabilities_enabled() + and _observability.is_otel_capabilities_enabled(client_options) ): try: from opentelemetry import trace @@ -293,7 +292,7 @@ def wrap_method( with_call=False, client_options=None, method_name=None, - trace=True, + is_streaming=False, ): """Wrap an RPC method with common behavior. @@ -383,10 +382,9 @@ def get_topic(name, timeout=None): custom OpenTelemetry tracer providers. Defaults to None. method_name (Optional[str]): Optional explicit full RPC method name (e.g. "/google.cloud.secretmanager.v1.SecretManagerService/AccessSecretVersion"). - Used to identify the RPC for observability and tracing. - trace (bool): Whether to create OpenTelemetry Tier 3 tracing spans for this - callable. Defaults to True. When False, or when method_name is None, - tracing spans are bypassed. + Used to identify the RPC for observability. + is_streaming (bool): Whether the RPC method is streaming. Defaults to False. + Streaming methods are currently gated and do not generate Tier 3 spans. Returns: Callable: A new callable that takes optional ``retry``, ``timeout``, @@ -416,6 +414,6 @@ def get_topic(name, timeout=None): metadata=user_agent_metadata, client_options=client_options, method_name=method_name, - trace=trace, + is_streaming=is_streaming, ) ) diff --git a/packages/google-api-core/google/api_core/gapic_v1/method_async.py b/packages/google-api-core/google/api_core/gapic_v1/method_async.py index 9ea110717352..b0e6c816cedb 100644 --- a/packages/google-api-core/google/api_core/gapic_v1/method_async.py +++ b/packages/google-api-core/google/api_core/gapic_v1/method_async.py @@ -40,7 +40,7 @@ def wrap_method( *, client_options=None, method_name=None, - trace=True, + is_streaming=False, ): """Wrap an async RPC method with common behavior. @@ -63,6 +63,6 @@ def wrap_method( metadata=metadata, client_options=client_options, method_name=method_name, - trace=trace, + is_streaming=is_streaming, ) ) diff --git a/packages/google-api-core/tests/unit/gapic/test_method.py b/packages/google-api-core/tests/unit/gapic/test_method.py index 5f3330c8a689..beef840caf7c 100644 --- a/packages/google-api-core/tests/unit/gapic/test_method.py +++ b/packages/google-api-core/tests/unit/gapic/test_method.py @@ -394,8 +394,8 @@ def test_wrap_method_otel_tracing_omitted_method_name_skips_span(monkeypatch): mock_trace.get_tracer.assert_not_called() -def test_wrap_method_otel_tracing_explicit_trace_false_skips_span(monkeypatch): - """Proves that when trace=False is explicitly passed (e.g. streaming call), no span is created.""" +def test_wrap_method_otel_tracing_streaming_skips_span(monkeypatch): + """Proves that when is_streaming=True is passed, no Tier 3 span is created.""" monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") mock_target = mock.Mock(return_value="success") @@ -417,7 +417,7 @@ def test_wrap_method_otel_tracing_explicit_trace_false_skips_span(monkeypatch): wrapped = google.api_core.gapic_v1.method.wrap_method( mock_target, method_name="/google.cloud.secretmanager.v1.SecretManagerService/StreamingRead", - trace=False, + is_streaming=True, ) result = wrapped() @@ -685,8 +685,8 @@ def test_wrap_method_async_otel_tracing(monkeypatch): ) -def test_wrap_method_async_otel_tracing_trace_false_skips_span(monkeypatch): - """Proves that method_async.wrap_method with trace=False skips span creation.""" +def test_wrap_method_async_otel_tracing_streaming_skips_span(monkeypatch): + """Proves that method_async.wrap_method with is_streaming=True skips span creation.""" from google.api_core.gapic_v1 import method_async monkeypatch.setenv("GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED", "true") @@ -710,7 +710,7 @@ def test_wrap_method_async_otel_tracing_trace_false_skips_span(monkeypatch): mock_target, kind=None, method_name="google.test.AsyncService/AsyncMethod", - trace=False, + is_streaming=True, ) result = wrapped()