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..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 @@ -18,11 +18,12 @@ compression, pagination, and long-running operations to gRPC methods. """ +import contextlib import enum 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 @@ -104,6 +105,24 @@ def _extract_metrics_header(metadata) -> Tuple[str, List[Tuple[str, str]]]: return metric_str, arbitrary_metadata +def _extract_rpc_identity( + method_name: str, +) -> Tuple[str, str, str]: + """Extract (full_rpc_name, service_name, rpc_method_name) from an explicit method name. + + Args: + 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 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): """Callable that applies retry, timeout, and metadata logic. @@ -123,6 +142,16 @@ 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"). + 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__( @@ -132,11 +161,18 @@ def __init__( timeout, compression, metadata=None, + client_options=None, + method_name=None, + trace=True, ): self._target = target self._retry = retry self._timeout = timeout 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) self._static_metadata = tuple(remaining) @@ -148,6 +184,42 @@ def __init__( else: 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. + self._tracer = None + self._span_name = None + self._span_attributes = None + if ( + trace + and 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_service, self._rpc_method = ( + _extract_rpc_identity(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 ): @@ -186,7 +258,29 @@ def __call__( if self._compression is not None: kwargs["compression"] = compression - return wrapped_func(*args, **kwargs) + span_context_manager = contextlib.nullcontext() + if self._tracer is not None and self._span_name is not None: + try: + from opentelemetry import trace + + span_context_manager = self._tracer.start_as_current_span( + self._span_name, + kind=trace.SpanKind.CLIENT, + 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( @@ -197,6 +291,9 @@ def wrap_method( client_info=client_info.DEFAULT_CLIENT_INFO, *, with_call=False, + client_options=None, + method_name=None, + trace=True, ): """Wrap an RPC method with common behavior. @@ -280,6 +377,16 @@ 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"). + 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``, @@ -307,5 +414,8 @@ def get_topic(name, timeout=None): default_timeout, default_compression, 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 d361bf9f961f..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 @@ -37,6 +37,10 @@ def wrap_method( default_compression=None, client_info=client_info.DEFAULT_CLIENT_INFO, kind=_DEFAULT_ASYNC_TRANSPORT_KIND, + *, + client_options=None, + method_name=None, + trace=True, ): """Wrap an async RPC method with common behavior. @@ -57,5 +61,8 @@ def wrap_method( default_timeout, default_compression, metadata=metadata, + client_options=client_options, + method_name=method_name, + trace=trace, ) ) 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 fbe7f2a5f0f1..5f3330c8a689 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 @@ -26,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 @@ -346,3 +348,371 @@ 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") + + 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, + method_name="google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", + ) + assert wrapped() == "success" + mock_target.assert_called_once() + + +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_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) + result = wrapped() + + assert result == "success" + 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") + 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/ListSecrets", + ) + 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_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") + + 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" + + client_options = client_options_lib.ClientOptions(tracer_provider=mock_provider) + + 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, + client_options=client_options, + method_name="google.test.Service/TestMethod", + ) + 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_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_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" + + 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/ListSecrets", + ) + 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)) + + +def test_wrap_method_otel_tracing_bytes_method(monkeypatch): + """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_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, + default_timeout=60, + method_name=b"/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/ListSecrets", + kind="CLIENT", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.cloud.secretmanager.v1.SecretManagerService", + "rpc.method": "ListSecrets", + }, + ) + + +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") + + with ( + mock.patch( + "google.api_core._observability.is_otel_capabilities_enabled", + return_value=True, + ), + mock.patch.dict( + sys.modules, + { + "opentelemetry": None, + "opentelemetry.trace": None, + }, + ), + ): + wrapped = google.api_core.gapic_v1.method.wrap_method( + mock_target, + method_name="google.cloud.secretmanager.v1.SecretManagerService/ListSecrets", + ) + 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 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") + 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_provider = mock.Mock() + mock_provider.get_tracer.return_value = mock_tracer + + 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", + 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, + 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", + attributes={ + "rpc.system": "grpc", + "rpc.service": "google.test.AsyncService", + "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()