Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/google/api_core/_observability.py: 19%

Shortcuts on this page

r m x   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

100 statements  

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 )