1# -*- coding: utf-8 -*-
2# Copyright 2026 Google LLC
3#
4# Licensed under the Apache License, Version 2.0 (the "License");
5# you may not use this file except in compliance with the License.
6# You may obtain a copy of the License at
7#
8# http://www.apache.org/licenses/LICENSE-2.0
9#
10# Unless required by applicable law or agreed to in writing, software
11# distributed under the License is distributed on an "AS IS" BASIS,
12# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13# See the License for the specific language governing permissions and
14# limitations under the License.
15#
16
17"""OpenTelemetry helpers for resolving and instantiating interceptors."""
18
19from __future__ import annotations
20
21import urllib.parse
22from typing import TYPE_CHECKING, Any, Callable, Sequence
23
24from google.api_core import _feature_gating_helpers
25from google.api_core.client_options import ClientOptions
26
27if TYPE_CHECKING:
28 # flake8: grpc, trace, and ClientInterceptor are imported only for static analysis and type annotations
29 # The 'noqa: F401' comment avoids flake8 "imported but not used" errors.
30 import grpc # noqa: F401
31 import opentelemetry.trace # noqa: F401
32
33 from google.api_core.grpc_helpers import ClientInterceptor # noqa: F401
34
35_TRACER_PROVIDER = "tracer_provider"
36
37
38def is_otel_capabilities_enabled(
39 client_options: ClientOptions | dict[str, Any] | None = None,
40 env_var: str = "GOOGLE_SDK_EXPERIMENTAL_PYTHON_TRACING_ENABLED",
41) -> bool:
42 """Checks if OTel capabilities are enabled and installed.
43
44 Args:
45 client_options: The client options object or dictionary.
46 env_var: The environment variable to check for enablement.
47
48 Returns:
49 bool: True if enabled and installed, False otherwise.
50 """
51 is_tracing_enabled = _feature_gating_helpers.resolve_feature_flags(
52 env_var=env_var,
53 feature_key=_TRACER_PROVIDER,
54 configuration=client_options,
55 )
56
57 if is_tracing_enabled:
58 try:
59 import opentelemetry.instrumentation.grpc as otel_grpc # type: ignore[import-not-found] # noqa: F401
60
61 return True
62 except ImportError:
63 pass
64
65 return False
66
67
68def _extract_endpoint_attributes(
69 client_options: ClientOptions | dict[str, Any] | None = None,
70) -> dict[str, Any]:
71 """Extracts server.address, server.port (if non-default), and url.domain from client options if present.
72
73 Args:
74 client_options: The client options object or dictionary.
75
76 Returns:
77 dict[str, Any]: A dictionary containing url.domain and, if an api_endpoint is configured,
78 server.address and non-default server.port.
79 """
80 attrs: dict[str, Any] = {}
81 endpoint = None
82 universe_domain = None
83
84 if isinstance(client_options, dict):
85 endpoint = client_options.get("api_endpoint")
86 universe_domain = client_options.get("universe_domain")
87 elif client_options is not None:
88 endpoint = getattr(client_options, "api_endpoint", None)
89 universe_domain = getattr(client_options, "universe_domain", None)
90
91 attrs["url.domain"] = universe_domain or "googleapis.com"
92
93 if endpoint and isinstance(endpoint, str):
94 target = endpoint if "//" in endpoint else f"//{endpoint}"
95 parsed = None
96 hostname = None
97 port = None
98 try:
99 parsed = urllib.parse.urlsplit(target)
100 hostname = parsed.hostname
101 port = parsed.port
102 except ValueError:
103 pass
104
105 if hostname:
106 attrs["server.address"] = hostname
107 if port and parsed:
108 scheme = parsed.scheme.lower()
109 is_default_port = (port == 443 and scheme in ("https", "")) or (
110 port == 80 and scheme == "http"
111 )
112 if not is_default_port:
113 attrs["server.port"] = port
114 return attrs
115
116
117def _make_grpc_client_request_hook(
118 endpoint_attrs: dict[str, Any] | None = None,
119) -> Callable[[Any, Any], None]:
120 """Creates an OpenTelemetry gRPC client request hook with optional endpoint attributes.
121
122 Args:
123 endpoint_attrs: Optional static endpoint attributes to attach to every span.
124
125 Returns:
126 Callable[[Any, Any], None]: The request hook callback.
127 """
128 static_attrs = dict(endpoint_attrs) if endpoint_attrs else {}
129
130 def client_request_hook(span: Any, request: Any) -> None:
131 if span is None or not getattr(span, "is_recording", lambda: True)():
132 return
133
134 # Upstream opentelemetry-instrumentation-grpc may format span names with a
135 # leading slash (e.g. "/package.Service/Method"). Normalize the span name
136 # and ensure rpc.method is always captured as the clean, fully-qualified name.
137 span_name = getattr(span, "name", None)
138 if isinstance(span_name, str) and span_name:
139 clean_method_name = span_name.lstrip("/")
140 if span_name.startswith("/") and hasattr(span, "update_name"):
141 span.update_name(clean_method_name)
142 span.set_attribute("rpc.method", clean_method_name)
143
144 attrs: dict[str, Any] = {
145 "rpc.system.name": "grpc",
146 }
147 if static_attrs:
148 attrs.update(static_attrs)
149 for key, value in attrs.items():
150 span.set_attribute(key, value)
151
152 return client_request_hook
153
154
155_grpc_client_request_hook = _make_grpc_client_request_hook()
156
157
158def _grpc_client_response_hook(span: Any, response: Any) -> None:
159 """OpenTelemetry gRPC client response hook to record successful response status.
160
161 Upstream ``opentelemetry-instrumentation-grpc`` sets the integer status code
162 ``rpc.grpc.status_code`` (e.g. 0), but does not record the modern string status
163 ``rpc.response.status_code`` (e.g. "OK") required by Cloud Trace and current
164 OpenTelemetry semantic conventions (v1.27.0+).
165
166 This hook enriches successful RPC attempt spans with ``rpc.response.status_code = "OK"``.
167 Errors and non-OK statuses are handled at the Tier 3 method span layer or upstream.
168
169 Upstream handles synchronous and asynchronous invocations differently:
170 - **Synchronous gRPC**: Upstream only invokes the response hook when an RPC call
171 succeeds. On failure, the hook is bypassed entirely.
172 - **Asynchronous gRPC**: Upstream invokes the response hook unconditionally for
173 both successes and failures (passing exception details on error). However, it
174 always marks ``span.status`` with an error status before calling the hook.
175
176 Because of this disparity, this hook checks ``span.status`` to guard against
177 async failure callbacks while allowing synchronous and successful asynchronous
178 calls to be marked "OK".
179
180 Note:
181 If upstream ``opentelemetry-instrumentation-grpc`` adds native support for
182 modern ``rpc.response.status_code`` in future releases, this hook can be retired.
183
184 Args:
185 span: The OpenTelemetry span.
186 response: The gRPC response object or details.
187 """
188 if not span.is_recording():
189 return
190
191 # Guard against upstream async calls that invoke this hook on failures.
192 status = getattr(span, "status", None)
193 status_code = getattr(status, "status_code", None)
194 if (
195 getattr(status_code, "name", None) == "ERROR"
196 or getattr(status_code, "value", None) == 2
197 ):
198 return
199
200 span.set_attribute("rpc.response.status_code", "OK")
201
202
203def _get_tracer_provider(
204 client_options: ClientOptions | dict[str, Any] | None = None,
205) -> opentelemetry.trace.TracerProvider | None:
206 """Extracts the OpenTelemetry tracer provider from client options if present.
207
208 Args:
209 client_options: The client options object or dictionary.
210
211 Returns:
212 opentelemetry.trace.TracerProvider | None: The tracer provider if present,
213 None otherwise.
214 """
215 if isinstance(client_options, dict):
216 return client_options.get(_TRACER_PROVIDER)
217 elif client_options is not None:
218 return getattr(client_options, _TRACER_PROVIDER, None)
219 return None
220
221
222def get_otel_interceptor(
223 client_options: ClientOptions | dict[str, Any] | None = None,
224) -> Callable[[grpc.Channel], grpc.Channel] | None:
225 """Returns an interceptor callable that wraps a sync gRPC channel with OpenTelemetry tracing.
226
227 Args:
228 client_options: The client options object or dictionary used for feature gating
229 and extracting the tracer provider.
230
231 Returns:
232 Callable[[grpc.Channel], grpc.Channel] | None: An interceptor callable if OpenTelemetry
233 tracing is enabled and installed, None otherwise.
234 """
235 if not is_otel_capabilities_enabled(client_options):
236 return None
237
238 import opentelemetry.instrumentation.grpc as otel_grpc # type: ignore[import-not-found]
239
240 endpoint_attrs = _extract_endpoint_attributes(client_options)
241 request_hook = _make_grpc_client_request_hook(endpoint_attrs)
242
243 interceptor: ClientInterceptor = otel_grpc.client_interceptor(
244 tracer_provider=_get_tracer_provider(client_options),
245 request_hook=request_hook,
246 response_hook=_grpc_client_response_hook,
247 )
248
249 def otel_interceptor(channel: grpc.Channel) -> grpc.Channel:
250 return otel_grpc.intercept_channel(channel, interceptor)
251
252 return otel_interceptor
253
254
255def get_otel_async_interceptor(
256 client_options: ClientOptions | dict[str, Any] | None = None,
257) -> Sequence[grpc.aio.ClientInterceptor] | None:
258 """Returns async gRPC client interceptors for OpenTelemetry tracing.
259
260 Args:
261 client_options: The client options object or dictionary used for feature gating
262 and extracting the tracer provider.
263
264 Returns:
265 Sequence[grpc.aio.ClientInterceptor] | None: Instantiated OpenTelemetry async
266 client interceptors if tracing is enabled and installed, None otherwise.
267 """
268 if not is_otel_capabilities_enabled(client_options):
269 return None
270
271 # Ignored by mypy: Optional dependency only loaded if early-return is skipped
272 import opentelemetry.instrumentation.grpc as otel_grpc # type: ignore[import-not-found]
273
274 endpoint_attrs = _extract_endpoint_attributes(client_options)
275 request_hook = _make_grpc_client_request_hook(endpoint_attrs)
276
277 return otel_grpc.aio_client_interceptors(
278 tracer_provider=_get_tracer_provider(client_options),
279 request_hook=request_hook,
280 response_hook=_grpc_client_response_hook,
281 )